dbaai 发表于 7 天前

ETL 增量同步实战:时间戳轮询与 CDC 的取舍

ETL 增量同步实战:时间戳轮询与 CDC 的取舍

数据同步这件事,很多团队的起点都是一样的:写个定时任务,每天凌晨把业务库的表全量重抽一遍。数据量小的时候毫无问题,等表涨到几千万行、单表几十 GB,凌晨两点的窗口就再也装不下了。这时候所有人都会说"改成增量吧",但增量怎么做,才是真正见功力的地方。

一、具体的问题

先说一个真实的场景。订单表 t_order 约 4200 万行,每天新增约 12 万、更新约 30 万(状态流转)。原来的做法是每天 02:00 全量抽取,耗时 38 分钟,抽完再跑下游汇总,整条链路 06:30 才结束,业务部门 08:00 看报表,中间只剩一个半小时的缓冲。某天主库做了一次大版本发布,重试一次,报表直接延期到中午。

改成增量之后,第一个版本是"按 updated_at 取昨天之后的数据"。上线当天就出了三个问题:

问题一:删除的数据同步不过去。 业务侧做了订单合并,物理删掉了 3000 多行。源表没了,目标表还在,下游按订单号 sum 出来的金额比业务库多了 87 万。时间戳方案天然只能捕获 insert 和 update,delete 对它来说是隐形的。

问题二:边界数据漏抽。 抽取条件是 updated_at >= '2026-09-15 02:00:00',但有一笔事务在 01:59:58 开始、02:00:03 才提交。MySQL 的 updated_at 用的是事务开始时的 NOW()(或者说是语句执行时间),水位线切下去之后,这笔数据两边都够不着——昨天的批次查不到它(昨天跑的时候还没提交),今天的批次也查不到它(updated_at 小于今天的水位)。

问题三:字段没被维护。 有张表是老系统遗留的,updated_at 字段只有部分代码路径会更新,直接改表状态的那段代码压根没碰这个字段。结果这批更新全丢了。

这三个问题不是个例,是时间戳轮询方案的固有缺陷。要解决它们,得先搞清楚增量同步的几种"水位"到底是什么。

二、核心原理

1. 四种常见的增量水位


方式依据能抓 delete对源库压力侵入性
时间戳轮询updated_at 字段否中(需索引扫描)需字段被正确维护
自增主键id > last_max_id否低无
触发器影子表DML 触发器写日志表能高(事务内多写一次)强,DDL 要同步改
日志 CDCbinlog / WAL能极低无


自增主键只能覆盖"只增不改"的流水表,比如日志、埋点。触发器方案在生产上基本已经被淘汰了——它把同步逻辑塞进了业务事务里,触发器一旦报错,业务写就跟着失败。真正摆在选择台面上的,就是时间戳轮询和 CDC 两种。

2. 时间戳轮询为什么必须加"重叠窗口"

因为存在"长事务"和"时钟抖动"。正确做法不是 updated_at >= 上次水位,而是:


WHERE updated_at >= (上次水位 - 重叠窗口)
AND updated_at <(本次水位)


重叠窗口一般取 5~15 分钟,取值依据是源库上最长事务的耗时。这样会把一部分数据重复抽一遍,所以下游必须是幂等的 upsert,不能是 insert。接受"至少一次投递 + 幂等覆盖",是设计增量链路的前提,不要在这里纠结。

3. CDC 为什么能抓到 delete

CDC 读的是数据库的变更日志。MySQL 的 binlog(row 格式)里记录的是行级前后镜像:insert 记录新值,update 记录旧值和新值,delete 记录旧值。所以 delete 事件是完整可捕获的。它还有两个额外好处:一是不需要在源表上加索引去扫(对源库几乎没有额外压力),二是延迟可以做到秒级。

代价是:需要开 row 格式 binlog、需要一个消费端(Canal / Debezium / Flink CDC)、DDL 变更需要单独处理、并且同样要保证幂等。

4. 什么时候该上 CDC

我的判断标准很简单,满足任意一条就上 CDC:有物理删除、要求分钟级延迟、源表没有可用的更新时间字段、源库 CPU 已经吃紧扛不动轮询扫描。如果表只是追加写入、允许 T+1、且 updated_at 维护得很干净,时间戳轮询完全够用,没必要为了技术先进上 CDC——多一个组件就多一个故障点。

三、实例参考(动手步骤)

下面以 MySQL 8.0 源库为例,两种方案都给出可直接照做的操作。

步骤 1:建表并造测试数据


-- 源库业务表
CREATE TABLE t_order (
id          BIGINT PRIMARY KEY AUTO_INCREMENT,
order_no    VARCHAR(32) NOT NULL,
amount      DECIMAL(12,2) NOT NULL,
status      TINYINT NOT NULL DEFAULT 0,
updated_atDATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP
                     ON UPDATE CURRENT_TIMESTAMP,
KEY idx_updated_at (updated_at)
) ENGINE=InnoDB;

-- 目标库(数仓 ODS 层)
CREATE TABLE ods_order (
id          BIGINT PRIMARY KEY,
order_no    VARCHAR(32),
amount      DECIMAL(12,2),
status      TINYINT,
updated_atDATETIME,
etl_time    DATETIME,
is_deletedTINYINT DEFAULT 0
) ENGINE=InnoDB;

-- 水位表
CREATE TABLE etl_watermark (
task_name   VARCHAR(64) PRIMARY KEY,
last_valueDATETIME,
updated_atDATETIME
) ENGINE=InnoDB;


步骤 2:方案 A —— 时间戳轮询(带重叠窗口)


-- 1) 取上次水位
SELECT last_value FROM etl_watermark WHERE task_name = 'sync_order';

-- 2) 带 10 分钟重叠窗口抽取(假设上次水位 2026-09-15 02:00:00)
SELECT id, order_no, amount, status, updated_at
FROM t_order
WHERE updated_at >= '2026-09-15 01:50:00'
AND updated_at <'2026-09-16 02:00:00';


捞出来的数据做幂等 upsert,这一步不能用 insert:


INSERT INTO ods_order (id, order_no, amount, status, updated_at, etl_time)
VALUES (1001, 'NO20260915001', 299.00, 2, '2026-09-15 09:12:33', NOW())
ON DUPLICATE KEY UPDATE
order_no   = VALUES(order_no),
amount   = VALUES(amount),
status   = VALUES(status),
updated_at = VALUES(updated_at),
etl_time   = NOW();


最后推进水位(注意:水位取的是本次批次的上界,不是 MAX(updated_at)):


INSERT INTO etl_watermark (task_name, last_value, updated_at)
VALUES ('sync_order', '2026-09-16 02:00:00', NOW())
ON DUPLICATE KEY UPDATE last_value = VALUES(last_value), updated_at = NOW();


物理删除的兜底做法(如果暂时不上 CDC):每天一次全量主键对账,补标记删除。


-- 目标表有、源表没有的,标记删除
UPDATE ods_order o
LEFT JOIN t_order s ON s.id = o.id
SET o.is_deleted = 1, o.etl_time = NOW()
WHERE s.id IS NULL AND o.is_deleted = 0;


步骤 3:方案 B —— CDC(Debezium 最小可用配置)

先确认源库 binlog 格式:


SHOW VARIABLES LIKE 'binlog_format';      -- 必须是 ROW
SHOW VARIABLES LIKE 'binlog_row_image';   -- 建议 FULL
SHOW VARIABLES LIKE 'server_id';          -- 必须非 0 且唯一


不对就改配置文件(改完重启实例):



server_id      = 1001
log_bin          = mysql-bin
binlog_format    = ROW
binlog_row_image = FULL


注册连接器(Kafka Connect REST):


{
"name": "order-connector",
"config": {
    "connector.class": "io.debezium.connector.mysql.MySqlConnector",
    "database.hostname": "10.0.0.11",
    "database.port": "3306",
    "database.user": "cdc",
    "database.password": "******",
    "database.server.id": "5401",
    "database.include.list": "bizdb",
    "table.include.list": "bizdb.t_order",
    "database.history.kafka.bootstrap.servers": "kafka:9092",
    "database.history.kafka.topic": "schema-changes.bizdb",
    "snapshot.mode": "when_needed"
}
}


消费端按 op 字段(c 新增 / u 更新 / d 删除 / r 快照读)分流处理,d 事件直接打删除标记:


-- 收到 op='d' 的事件时
UPDATE ods_order SET is_deleted = 1, etl_time = NOW()
WHERE id = 1001;


步骤 4:两种方案的效果对比

同一张 4200 万行的 t_order,改造前后的实测数据:


指标全量重抽时间戳轮询Debezium CDC
单次耗时38 min42 s准实时(<3 s)
单次传输行数4200 万约 42 万约 42 万(事件数)
能否捕获删除能否(需对账补)能
源库额外负载高(全表扫)中(索引范围扫)极低(读 binlog)
首次全量不需要需要(先做一次基线)需要(snapshot)
运维复杂度低低中高


我们最后的落地选择是:核心交易表(订单、支付、库存)走 CDC,配置类和字典类小表继续时间戳轮询,日志类流水表走自增主键。不是所有表都值得上 CDC,按表的变更特征分档,成本才压得住。

四、实操检查清单


[*][ ] 梳理表的变更特征:只增 / 有更新 / 有物理删除,按分档决定用哪种水位
[*][ ] 确认 updated_at 字段被所有写路径维护,且有索引 idx_updated_at
[*][ ] 确认源库时区与抽取程序时区一致,避免时间列偏移 8 小时
[*][ ] 时间字段用 DATETIME(6) 或 TIMESTAMP,避免同一毫秒内多行导致边界抖动
[*][ ] 轮询 SQL 必须带重叠窗口(5~15 分钟),窗口值 ≥ 源库最长事务耗时
[*][ ] 下游写入一律改成幂等 upsert,禁止裸 insert;主键或唯一键必须存在
[*][ ] 水位表的推进放在写入成功之后,且放在同一个事务里,避免"数据没写进去水位先涨了"
[*][ ] 水位取批次的查询上界,不要取 MAX(updated_at),防止空洞数据永久丢失
[*][ ] 有删除场景但暂不上 CDC 的,必须配每日主键全量对账,并给下游 is_deleted 过滤条件
[*][ ] 上 CDC 前检查 binlog_format=ROW、binlog_row_image=FULL、server_id 唯一
[*][ ] binlog 保留期 ≥ 最长可接受的中断时长(建议 ≥ 72 小时),否则断连后要重新做快照
[*][ ] CDC 的 server.id 不要和已有从库冲突,否则会导致复制拓扑混乱
[*][ ] DDL 变更单独走审批:加列可以自动兼容,改列名/改类型会导致消费端解析失败
[*][ ] 目标表保留 etl_time 与源端 updated_at 两个时间,问题排查时才能分清是谁的锅
[*][ ] 配置行数波动监控:本次抽取行数偏离近 7 日均值 ±50% 时告警拦停,别让它污染下游


增量同步这件事,技术上不复杂,复杂的是把"至少一次投递 + 幂等 + 对账"这三个东西当成一条完整链路来设计。只做其中一环,迟早会在某个月初对账的时候翻车。

—— dbaai
页: [1]
查看完整版本: ETL 增量同步实战:时间戳轮询与 CDC 的取舍