5 minutes
增量抽取与变更数据捕获(CDC)
为什么全量刷新不够用
在 ETL 入门阶段,全量刷新听起来很省事:每次把源表的所有数据删掉,重新拉一遍。对于小数据量表,这没什么问题。但当数据量增长到百万、千万甚至亿级时,全量刷新就成了一场噩梦。
来看几个数字。假设你有一个订单表,每天新增 5 万条记录,总计 2000 万行。全量读取 2000 万行可能耗费 30 分钟,但实际发生变化的数据只有 5 万行,占总量的 0.25%。这意味着你浪费了 99.75% 的传输量和计算量。
全量刷新还有几个更棘手的问题:
- 窗口时间不够:数据量越大,刷新窗口越长。如果每天只有 2 小时的维护窗口,全量同步 10 亿条数据根本跑不完
- 对源库压力大:全表扫描会对业务数据库造成巨大压力,影响在线交易
- 浪费存储和带宽:每次传输全部数据,压缩、网络、磁盘 IO 的成本成倍增加
增量抽取就是为了解决这些问题而生的。
增量抽取的三种主要策略
增量抽取的核心是回答一个问题:“上次处理完之后,哪些数据发生了变化?” 针对不同的数据源和业务场景,有三种主流策略。
时间戳增量
时间戳策略是最简单、应用最广泛的增量抽取方式。原理是:源表里有一个时间字段(如 updated_at、modified_date),记录每条数据最后修改的时间。每次抽取时,只取时间大于上次抽取时间的记录。
-- 最后一次抽取的时间是 2026-04-12 03:00:00
-- 只抽取此后修改的数据
SELECT * FROM orders
WHERE updated_at > '2026-04-12 03:00:00'
AND updated_at <= '2026-04-13 03:00:00';
时间戳策略的优点是实现简单,对源库改造小。但缺点也很明显:
- 需要源表有时间戳字段,且业务代码正确维护该字段
- 物理删除的数据无法被捕获(删除不更新时间戳)
- 时间精度不够时可能漏数据(如果两行在同一毫秒内修改)
- 批量更新时间戳的情况可能导致重复或遗漏
增量文件/日志
有些系统会主动生成包含变更的文件或日志。例如:
- 应用日志:业务系统把每次增删改操作输出到日志文件
- 增量文件:ERP 系统每天生成变更文件,记录当天修改的记录 ID
- Binlog:MySQL 的二进制日志记录了所有数据变更
增量文件的方式对源系统侵入最小,但依赖应用配合生成这些文件。
基于日志的 CDC(Change Data Capture)
CDC 是增量抽取的终极方案。它的原理是直接读取数据库的事务日志(如 MySQL 的 binlog、PostgreSQL 的 WAL),从中解析出所有数据变更事件。
-- 查看 MySQL 是否开启 binlog
SHOW VARIABLES LIKE 'log_bin';
SHOW VARIABLES LIKE 'binlog_format';
-- 确保 binlog 格式为 ROW
-- ROW 格式记录每行数据变更前后的值,是最适合 CDC 的模式
基于日志的 CDC 有几个压倒性的优势:
- 零侵入:不需要修改业务代码,不需要在表上加时间戳
- 捕获所有变更:包括 INSERT、UPDATE、DELETE,连 DDL 都能捕获
- 近乎实时:数据一旦提交,CDC 工具马上就能感知到
- 一致性强:直接从事务日志读,保证不丢数据
CDC 三种方法详解
从实现层面看,CDC 有三种具体的方法,我们逐个分析。
数据库触发器 CDC
在源数据库上建立触发器,每当有 INSERT、UPDATE、DELETE 操作时,触发器把变更记录写入一张专门的审计表(Changelog Table)。
-- 创建变更记录表
CREATE TABLE orders_changelog (
log_id BIGINT AUTO_INCREMENT PRIMARY KEY,
table_name VARCHAR(100),
operation_type ENUM('INSERT', 'UPDATE', 'DELETE'),
record_id BIGINT,
old_data JSON,
new_data JSON,
changed_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
processed BOOLEAN DEFAULT FALSE
);
-- 创建触发器捕获 INSERT
CREATE TRIGGER trg_orders_insert
AFTER INSERT ON orders
FOR EACH ROW
INSERT INTO orders_changelog (table_name, operation_type, record_id, new_data)
VALUES ('orders', 'INSERT', NEW.id, JSON_OBJECT(
'order_id', NEW.order_id,
'amount', NEW.amount,
'status', NEW.status
));
-- 创建触发器捕获 UPDATE
CREATE TRIGGER trg_orders_update
AFTER UPDATE ON orders
FOR EACH ROW
INSERT INTO orders_changelog (table_name, operation_type, record_id, old_data, new_data)
VALUES ('orders', 'UPDATE', NEW.id,
JSON_OBJECT(
'order_id', OLD.order_id,
'amount', OLD.amount,
'status', OLD.status
),
JSON_OBJECT(
'order_id', NEW.order_id,
'amount', NEW.amount,
'status', NEW.status
)
);
触发器方案的优点是全兼容,任何支持触发器的数据库都能用。缺点是对源库性能有影响——每次 DML 操作都额外多写一次变更表。在高并发场景下,触发器可能成为瓶颈。
事务日志挖掘 CDC(Debezium)
这是目前最主流的 CDC 方式。Debezium 是一个基于 Kafka Connect 的开源 CDC 平台,支持 MySQL、PostgreSQL、MongoDB、SQL Server 等多种数据库。
Debezium 的工作原理如下:
- 连接到数据库,读取 binlog 或 WAL 日志
- 解析日志中的变更事件
- 把变更事件写入 Kafka 主题
- 下游消费者从 Kafka 读取并处理变更
一个 Debezium 捕获的变更事件结构如下:
{
"schema": { "...": "..." },
"payload": {
"op": "u",
"ts_ms": 1712880000123,
"before": {
"order_id": 1001,
"amount": 299.00,
"status": "PAID"
},
"after": {
"order_id": 1001,
"amount": 299.00,
"status": "SHIPPED"
},
"source": {
"db": "ecommerce",
"table": "orders",
"ts_ms": 1712880000100
}
}
}
在 Debezium 中,每条变更记录包含操作类型(op)、变更前的数据(before)、变更后的数据(after)以及源信息(source)。其中 op 字段的含义是:
c:创建(Create),对应 INSERTu:更新(Update),对应 UPDATEd:删除(Delete),对应 DELETEr:读取(Read),用于快照
使用 Docker 启动 Debezium 连接 MySQL 的步骤非常简单:
version: '3.8'
services:
zookeeper:
image: confluentinc/cp-zookeeper:7.6.0
environment:
ZOOKEEPER_CLIENT_PORT: 2181
ZOOKEEPER_TICK_TIME: 2000
kafka:
image: confluentinc/cp-kafka:7.6.0
ports:
- "9092:9092"
depends_on:
- zookeeper
environment:
KAFKA_BROKER_ID: 1
KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092
KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
debezium-connector:
image: debezium/connect:2.5
ports:
- "8083:8083"
depends_on:
- kafka
environment:
BOOTSTRAP_SERVERS: kafka:9092
GROUP_ID: 1
CONFIG_STORAGE_TOPIC: connect-configs
OFFSET_STORAGE_TOPIC: connect-offsets
STATUS_STORAGE_TOPIC: connect-status
启动服务后,注册一个 Debezium 连接器:
curl -i -X POST localhost:8083/connectors \
-H "Content-Type: application/json" \
-d '{
"name": "orders-connector",
"config": {
"connector.class": "io.debezium.connector.mysql.MySqlConnector",
"database.hostname": "mysql-host",
"database.port": "3306",
"database.user": "debezium",
"database.password": "dbz_pass",
"database.server.id": "184054",
"database.server.name": "ecommerce",
"database.include.list": "ecommerce",
"table.include.list": "ecommerce.orders",
"database.history.kafka.bootstrap.servers": "kafka:9092",
"database.history.kafka.topic": "schema-changes.orders",
"include.schema.changes": "true",
"snapshot.mode": "initial"
}
}'
注册完成后,Debezium 会自动读取 MySQL 的 binlog,把 orders 表的每一次变更推送到 Kafka 主题 ecommerce.ecommerce.orders。其他系统只需要从这个主题消费就能获取实时的数据变更。
Debezium 是目前最流行的 CDC 工具之一。它的优势在于与 Kafka 深度集成,天然支持消息持久化和回溯消费。但缺点是运维复杂度较高,需要同时维护 Kafka 和 Kafka Connect 集群。
查询式 CDC
查询式 CDC 不做日志解析,而是通过 SQL 查询来找到变更数据。常见做法是在源表上维护水印字段,每次查询大于水印值的数据。
import psycopg2
from datetime import datetime
class QueryBasedCDC:
"""基于查询的 CDC 实现"""
def __init__(self, conn_params, watermark_table="etl_watermarks"):
self.conn = psycopg2.connect(**conn_params)
self.watermark_table = watermark_table
self._init_watermark_table()
def _init_watermark_table(self):
"""初始化水印表"""
with self.conn.cursor() as cur:
cur.execute(f"""
CREATE TABLE IF NOT EXISTS {self.watermark_table} (
source_table VARCHAR(255) PRIMARY KEY,
last_max_id BIGINT DEFAULT 0,
last_max_timestamp TIMESTAMP,
updated_at TIMESTAMP DEFAULT NOW()
)
""")
self.conn.commit()
def get_watermark(self, table_name):
"""获取水印值"""
with self.conn.cursor() as cur:
cur.execute(
f"SELECT last_max_id FROM {self.watermark_table} WHERE source_table = %s",
(table_name,)
)
row = cur.fetchone()
return row[0] if row else 0
def update_watermark(self, table_name, max_id):
"""更新水印值"""
with self.conn.cursor() as cur:
cur.execute(f"""
INSERT INTO {self.watermark_table} (source_table, last_max_id, updated_at)
VALUES (%s, %s, NOW())
ON CONFLICT (source_table)
DO UPDATE SET last_max_id = %s, updated_at = NOW()
""", (table_name, max_id, max_id))
self.conn.commit()
def fetch_incremental(self, source_table, batch_size=5000):
"""增量获取数据"""
last_id = self.get_watermark(source_table)
with self.conn.cursor() as cur:
cur.execute(f"""
SELECT * FROM {source_table}
WHERE id > %s
ORDER BY id
LIMIT %s
""", (last_id, batch_size))
columns = [desc[0] for desc in cur.description]
rows = cur.fetchall()
if not rows:
return []
max_id = rows[-1][0]
result = [dict(zip(columns, row)) for row in rows]
self.update_watermark(source_table, max_id)
return result
# 使用示例
cdc = QueryBasedCDC({
"host": "localhost",
"port": 5432,
"dbname": "ecommerce",
"user": "etl_user",
"password": "etl_pass"
})
while True:
batch = cdc.fetch_incremental("orders", batch_size=5000)
if not batch:
print("没有新数据,等待 30 秒...")
time.sleep(30)
continue
print(f"获取到 {len(batch)} 条新订单")
process_orders(batch)
查询式 CDC 的优点是实现简单,不需要额外的基础设施。缺点是延迟较高(依赖定时轮询),并且无法捕获物理删除。
CDC 工具的对比
在实际项目中,选择合适的 CDC 工具需要考虑多个维度。下表对比了三种常见的工具选型:
| 特性 | Debezium | AWS DMS | Fivetran |
|---|---|---|---|
| 部署方式 | 自托管 | 托管服务 | SaaS |
| 支持的源库 | MySQL, PG, SQL Server, MongoDB 等 | 20+ 种 | 150+ 种 |
| 目标端 | Kafka | 多种数据库/S3 | 数据仓库 |
| 延迟 | 秒级 | 秒级 | 分钟级 |
| 运维复杂度 | 高 | 低 | 无 |
| 成本 | 开源免费 | 按实例计费 | 按行数计费 |
| 数据转换 | 需自建 | 可配置映射 | 内置 schema 迁移 |
选型建议:
- 如果团队有 Kafka 运维能力,Debezium 是最灵活的选择
- 如果不想操心基础设施,AWS DMS 或 Fivetran 这类托管服务更省心
- 如果是简单的增量抽取场景,查询式 CDC + 水印表就足够了
处理延迟到达的数据
现实世界中的数据不是完美的。由于网络延迟、批处理排队、上游系统故障等原因,部分数据可能会延迟到达。具体来说:
- 某笔订单昨晚 10 点创建,但今早 8 点才到达 ETL 系统
- 上游系统补录了一条昨天的数据,更新时间戳是今天
处理延迟数据有几种策略:
基于事件时间的处理
不依赖数据到达时间,而是依赖数据本身的业务时间(如 order_created_at)。这种方式在 Flink 和 Spark Streaming 中有成熟的支持。
def process_with_event_time(record):
"""基于事件时间处理"""
# 使用业务时间而不是系统时间
business_date = record["order_created_at"].date()
target_partition = f"orders_{business_date.strftime('%Y%m%d')}"
write_to_partition(record, target_partition)
允许更新的水印策略
水印不应该是单向增大的。允许一定的时间窗口内水印回退或重新处理某个时间段的数据。
class SlidingWatermark:
"""滑动窗口水印管理"""
def __init__(self, lookback_hours=24):
self.watermark_table = "etl_watermarks"
self.lookback = timedelta(hours=lookback_hours)
def get_safe_window(self):
"""获取安全处理窗口"""
last_watermark = self.read_watermark()
window_start = last_watermark - self.lookback
return window_start, last_watermark
def process_late_data(self, late_record):
"""处理延迟到达的数据"""
event_time = late_record["event_time"]
safe_start, safe_end = self.get_safe_window()
if safe_start <= event_time <= safe_end:
# 在安全窗口内,直接写入
self.upsert_to_target(late_record)
elif event_time < safe_start:
# 太旧的数据,需要人工确认
self.send_to_dead_letter(late_record, "数据超出回溯窗口")
精确一次处理语义
在 CDC 场景中,最糟糕的结果不是数据重复,而是数据丢失。重复数据可以通过去重来纠正,丢失的数据则很难发现。
实现精确一次处理需要三个层面的保证:
源端:至少一次捕获
CDC 工具需要记录已经读取到的日志位置(offset)。常见做法是把 offset 存到 Kafka 或独立的元数据存储中。
def capture_with_offset(debezium_record, offset_storage):
"""带偏移量记录的 CDC 消费"""
offset = debezium_record["source"]["offset"]
# 检查是否处理过
if offset_storage.is_processed(offset):
print("这条变更已经处理过,跳过")
return
# 记录偏移量
offset_storage.mark_processing(offset)
try:
# 处理变更
process_change(debezium_record)
offset_storage.mark_completed(offset)
except Exception as e:
print(f"处理失败: {e}")
offset_storage.mark_failed(offset)
目标端:幂等写入
前面文章提过,幂等性是批处理的核心原则。在 CDC 场景中同样适用。最好的做法是使用 UPSERT(MERGE)来写入目标表。
def upsert_to_warehouse(records, target_table):
"""幂等写入目标表"""
# 使用 MERGE 语句避免重复
merge_sql = f"""
MERGE INTO {target_table} AS target
USING (SELECT %s AS id, %s AS amount, %s AS status) AS source
ON target.id = source.id
WHEN MATCHED THEN
UPDATE SET amount = source.amount, status = source.status
WHEN NOT MATCHED THEN
INSERT (id, amount, status) VALUES (source.id, source.amount, source.status)
"""
for record in records:
execute(merge_sql, (record["id"], record["amount"], record["status"]))
协调:两阶段提交或事务
对于要求极高一致性的场景,可以结合事务性消息和两阶段提交。Kafka 支持事务性写入,配合下游的事务性读取可以实现端到端精确一次。
# Kafka 事务性生产
producer.init_transactions()
producer.begin_transaction()
try:
for record in changes:
producer.send("cdc-orders", record)
producer.commit_transaction()
except Exception:
producer.abort_transaction()
小结
增量抽取和 CDC 是解决大规模数据同步问题的核心手段。从简单的时间戳策略到基于日志的 Debezium,每种方法都有自己的适用场景。选择 CDC 方案时要考虑数据量大小、实时性要求、源库类型和团队运维能力。延迟到达的数据要采用事件时间处理或滑动窗口来应对,而精确一次处理依赖源头 offset 记录、目标端幂等写入、以及全链路的容错设计。
Summary: 增量抽取的三种策略和 CDC 原理,涵盖时间戳、日志、查询式方案及精确一次处理。