数据抽取概述

数据抽取是 ETL 流程的第一个环节,也是后续所有工作的基础。如果把 ETL 比作烹饪,抽取就是"采购食材"——食材的质量和新鲜度直接决定了最终菜品的好坏。

在实际工作中,数据抽取面临的挑战往往比想象中大:源系统的类型五花八门、数据量可能非常庞大、抽取过程不能影响业务系统的正常运行、网络不稳定导致连接中断等。因此,设计一个健壮的数据抽取方案是整个 ETL 工程的起点。

数据源的类型

现代数据架构中的数据源可以大致归为以下几类:

关系型数据库

关系型数据库是最常见的数据源。企业核心业务数据通常存储在 MySQL、PostgreSQL、Oracle、SQL Server 等系统中。从关系型数据库抽取数据通常有几种方式:直接执行 SELECT 查询、使用数据库的导出工具、通过日志复制(CDC)等。

文件系统

文件是另一种常见的数据载体。CSV 和 JSON 是最通用的交换格式,几乎任何系统都能读写。在大数据生态中,Parquet 和 Avro 这类列式存储格式越来越流行,它们具有更好的压缩率和查询性能。

API 接口

许多 SaaS 平台(如 Salesforce、HubSpot、Google Analytics)只通过 API 暴露数据。RESTful API 是目前的主流协议,GraphQL 也在逐渐普及。API 抽取通常需要考虑认证、限流、分页等问题。

消息队列和流数据

Kafka、RabbitMQ、AWS Kinesis 等消息系统承载着实时数据流。从这类系统抽取数据时,通常使用消费者模式持续订阅数据。这类场景往往偏向实时或准实时处理。

数据源类型 典型例子 抽取方式 常见挑战
关系型数据库 MySQL, PostgreSQL, Oracle JDBC/ODBC 查询、导出 dump 大表全量扫描影响性能
文件 CSV, JSON, Parquet, Avro 文件读取、FTP/S3 下载 文件编码、格式不一致
API REST, GraphQL, SOAP HTTP 请求 限流、认证、分页
消息队列 Kafka, RabbitMQ, Kinesis Consumer 订阅 消息顺序、重复消费

从关系型数据库抽取数据

JDBC 方式(Python 示例)

JDBC(Java Database Connectivity)是连接数据库的标准接口。在 Python 中,我们可以使用数据库驱动库来实现类似的功能。

import psycopg2
import pandas as pd

# 连接 PostgreSQL
conn = psycopg2.connect(
    host="localhost",
    port=5432,
    database="production_db",
    user="etl_user",
    password="secure_password"
)

# 全量抽取一张表
query = "SELECT * FROM orders WHERE order_date >= '2026-01-01'"
df = pd.read_sql(query, conn)
print(f"抽取到 {len(df)} 条订单记录")

conn.close()

分批抽取避免大表压力

对于大表,一次性的全量查询可能锁住大量行,导致业务系统响应变慢。一个常用的优化策略是按主键或时间字段分批抽取。

import psycopg2
import pandas as pd

def extract_in_batches(conn, table_name, key_column, batch_size=10000):
    """按主键分批抽取数据"""
    # 获取主键的最小值和最大值
    cursor = conn.cursor()
    cursor.execute(f"SELECT MIN({key_column}), MAX({key_column}) FROM {table_name}")
    min_val, max_val = cursor.fetchone()
    cursor.close()

    all_data = []
    current_min = min_val

    while current_min <= max_val:
        current_max = current_min + batch_size - 1
        query = f"""
            SELECT * FROM {table_name}
            WHERE {key_column} BETWEEN {current_min} AND {current_max}
        """
        batch = pd.read_sql(query, conn)
        all_data.append(batch)
        print(f"抽取 {table_name}: 主键范围 {current_min} - {current_max}, 共 {len(batch)} 行")
        current_min = current_max + 1

    return pd.concat(all_data, ignore_index=True)

# 使用示例
conn = psycopg2.connect("dbname=production_db user=etl_user")
df = extract_in_batches(conn, "orders", "id", batch_size=50000)
print(f"累计抽取 {len(df)} 条记录")
conn.close()

使用数据库导出工具

对于超大规模的数据迁移,直接用 SELECT 查询往往不够高效。数据库自带的导出工具经过深度优化,速度更快。

# PostgreSQL 导出
psql -h localhost -U etl_user -d production_db \
  -c "\copy orders TO 'orders_export.csv' WITH CSV HEADER"

# MySQL 导出
mysql -h localhost -u etl_user -p production_db \
  -e "SELECT * FROM orders" \
  --batch --quick > orders_export.tsv

# 使用 pg_dump 导出为自定义格式
pg_dump -h localhost -U etl_user \
  -t orders --data-only --format=custom \
  -f orders_export.dump production_db

从文件系统抽取数据

CSV 文件

CSV 是最通用的数据交换格式,但也是最容易出问题的格式之一——编码不一致、分隔符不统一、引号处理错误等等。

import pandas as pd
import chardet

def extract_csv(file_path):
    """智能读取 CSV 文件,自动检测编码和分隔符"""
    # 检测文件编码
    with open(file_path, "rb") as f:
        raw = f.read(10000)
        encoding = chardet.detect(raw)["encoding"]
    print(f"检测到编码: {encoding}")

    # 尝试读取
    try:
        df = pd.read_csv(file_path, encoding=encoding)
    except UnicodeDecodeError:
        # 如果编码检测失败,尝试常见编码
        for enc in ["utf-8", "utf-8-sig", "gbk", "gb2312", "latin1"]:
            try:
                df = pd.read_csv(file_path, encoding=enc)
                print(f"使用编码 {enc} 成功读取")
                break
            except UnicodeDecodeError:
                continue

    print(f"CSV 文件读取完成,共 {len(df)} 行,{len(df.columns)} 列")
    return df

# 使用示例
sales_data = extract_csv("sales_2026_q1.csv")

JSON 和嵌套数据

JSON 的灵活性使其成为 API 和 NoSQL 数据库的常用格式。但嵌套结构给数据分析带来了挑战。

import json
import pandas as pd

def extract_json_with_flatten(file_path):
    """读取 JSON 文件并展开嵌套字段"""
    with open(file_path, "r", encoding="utf-8") as f:
        data = json.load(f)

    # 如果 JSON 是数组形式
    if isinstance(data, list):
        df = pd.json_normalize(data)
    # 如果 JSON 包含嵌套对象,使用 json_normalize 展开
    elif isinstance(data, dict) and "data" in data:
        df = pd.json_normalize(data["data"])
    else:
        df = pd.json_normalize([data])

    print(f"JSON 文件读取完成,共 {len(df)} 行")
    return df

# 使用示例
user_data = extract_json_with_flatten("users_export.json")

Parquet 文件

Parquet 是列式存储格式,在大数据生态中被广泛使用。它的压缩率高、查询性能好,尤其适合存储大规模分析数据。

import pandas as pd

def extract_parquet(file_path):
    """读取 Parquet 文件"""
    df = pd.read_parquet(file_path)
    print(f"Parquet 文件读取完成")
    print(f"行数: {len(df)}, 列数: {len(df.columns)}")
    print(f"文件大小: {df.memory_usage(deep=True).sum() / 1024 / 1024:.2f} MB")
    return df

# 使用示例
clickstream_data = extract_parquet("clickstream_events.parquet")

从 API 抽取数据

REST API 抽取

从 REST API 抽取数据时,通常需要处理认证、分页和限流三个核心问题。

import requests
import time
import pandas as pd

def extract_from_api(base_url, api_key, page_size=100):
    """从 REST API 抽取分页数据"""
    headers = {"Authorization": f"Bearer {api_key}"}
    all_records = []
    page = 1

    while True:
        params = {
            "page": page,
            "per_page": page_size
        }

        response = requests.get(base_url, headers=headers, params=params)

        if response.status_code == 429:
            # 触发限流,等待后重试
            retry_after = int(response.headers.get("Retry-After", 60))
            print(f"触发限流,等待 {retry_after} 秒...")
            time.sleep(retry_after)
            continue

        response.raise_for_status()
        data = response.json()

        # 假设 API 返回 { "data": [...], "total_pages": N }
        all_records.extend(data["data"])
        total_pages = data.get("total_pages", 1)

        print(f"已抽取第 {page}/{total_pages} 页")
        page += 1

        if page > total_pages:
            break

        # 控制请求频率,避免触发限流
        time.sleep(0.5)

    return pd.DataFrame(all_records)

# 使用示例
# df = extract_from_api("https://api.example.com/v1/orders", "your-api-key")

变更数据捕获(CDC)

CDC 是一种增量抽取技术,它通过读取数据库的事务日志(WAL、binlog)来捕获数据变更,而不是定期全量扫描表。相比全量导入+增量更新的方式,CDC 的效率和实时性都要高得多。

Debezium + Kafka 实现 CDC

Debezium 是一个开源的 CDC 平台,基于 Kafka Connect 架构,支持 MySQL、PostgreSQL、MongoDB 等多种数据库。

# 使用 Docker 启动 Debezium 连接器
docker run -it --rm --name connect \
  -p 8083:8083 \
  -e BOOTSTRAP_SERVERS=kafka:9092 \
  -e GROUP_ID=1 \
  -e CONFIG_STORAGE_TOPIC=connect_configs \
  -e OFFSET_STORAGE_TOPIC=connect_offsets \
  debezium/connect:2.5
# 注册 MySQL 连接器(通过 REST API)
curl -X POST -H "Content-Type: application/json" \
  --data '{
    "name": "mysql-connector",
    "config": {
      "connector.class": "io.debezium.connector.mysql.MySqlConnector",
      "database.hostname": "mysql-host",
      "database.port": 3306,
      "database.user": "debezium",
      "database.password": "dbz_password",
      "database.server.id": 184054,
      "topic.prefix": "fulfillment",
      "database.include.list": "order_db",
      "table.include.list": "order_db.orders"
    }
  }' \
  http://localhost:8083/connectors

全量抽取 vs 增量抽取

维度 全量抽取 增量抽取
数据量 每次读取全部数据 只读取变化的数据
对源系统影响 大(可能锁表) 小
实现复杂度 简单 较复杂(需要维护偏移量)
数据一致性 快照一致性 取决于实现方式
适用场景 小表、首次同步 大表、频繁同步
运行频率 较低(每天一次) 较高(分钟级甚至实时)

抽取层的错误处理

数据抽取环节容易出现各种问题。一个健壮的抽取流程需要包含完善的错误处理机制。

import logging
from datetime import datetime
from functools import wraps

logging.basicConfig(level=logging.INFO)
logger = logging.getLogger("extract")

def retry_on_failure(max_retries=3, delay=5):
    """抽取重试装饰器"""
    def decorator(func):
        @wraps(func)
        def wrapper(*args, **kwargs):
            for attempt in range(1, max_retries + 1):
                try:
                    return func(*args, **kwargs)
                except Exception as e:
                    logger.warning(f"第 {attempt} 次抽取失败: {e}")
                    if attempt == max_retries:
                        logger.error(f"重试 {max_retries} 次后仍然失败")
                        raise
                    time.sleep(delay * attempt)  # 指数退避
            return None
        return wrapper
    return decorator

@retry_on_failure(max_retries=3, delay=5)
def extract_orders(date):
    """抽取订单数据,失败自动重试"""
    conn = psycopg2.connect("dbname=production_db")
    try:
        query = "SELECT * FROM orders WHERE order_date = %s"
        df = pd.read_sql(query, conn, params=[date])
        logger.info(f"成功抽取 {date} 的订单数据,共 {len(df)} 条")
        return df
    finally:
        conn.close()

小结

数据抽取是 ETL 流程的起点,也是决定整个流程质量的关键环节。本文介绍了关系型数据库、文件、API 和消息队列四种主流数据源的抽取方式,并通过 Python 代码展示了每种方式的实现细节。我们还讨论了全量和增量两种抽取策略的优劣,以及 CDC 技术在增量抽取中的重要作用。记住,一个好的抽取方案不仅要能拿到数据,还要保证不破坏源系统的稳定性,并且能够在失败时优雅恢复。

下一篇文章我们将进入 ETL 最核心的步骤——数据转换,学习如何对原始数据进行清洗、加工和重组。

Summary: 数据源类型与抽取方式,JDBC/文件/API 抽取示例及 CDC 简介。