4 minutes
数据抽取(Extract):数据源与连接
数据抽取概述
数据抽取是 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 简介。