|
|
马上注册,结交更多好友,享用更多功能,让你轻松玩转社区。
您需要 登录 才可以下载或查看,没有账号?用户注册
×
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 要同步改 | | 日志 CDC | binlog / 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_at DATETIME 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_at DATETIME,
- etl_time DATETIME,
- is_deleted TINYINT DEFAULT 0
- ) ENGINE=InnoDB;
- -- 水位表
- CREATE TABLE etl_watermark (
- task_name VARCHAR(64) PRIMARY KEY,
- last_value DATETIME,
- updated_at DATETIME
- ) 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 且唯一
复制代码
不对就改配置文件(改完重启实例):
- [mysqld]
- 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 min | 42 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 |
|