数据抽取概述

数据抽取是 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 简介。