前事不忘,后事之师,不忘国耻!

 用户注册  找回密码
 用户注册
搜索
查看: 9|回复: 0

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

[复制链接]

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

[复制链接]
dbaai

主题

0

回帖

181

积分

DBAAI

积分
181
14 小时前 | 显示全部楼层 |阅读模式

马上注册,结交更多好友,享用更多功能,让你轻松玩转社区。

您需要 登录 才可以下载或查看,没有账号?用户注册

×
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 >= 上次水位,而是:
  1. WHERE updated_at >= (上次水位 - 重叠窗口)
  2.   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:建表并造测试数据
  1. -- 源库业务表
  2. CREATE TABLE t_order (
  3.   id          BIGINT PRIMARY KEY AUTO_INCREMENT,
  4.   order_no    VARCHAR(32) NOT NULL,
  5.   amount      DECIMAL(12,2) NOT NULL,
  6.   status      TINYINT NOT NULL DEFAULT 0,
  7.   updated_at  DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP
  8.                        ON UPDATE CURRENT_TIMESTAMP,
  9.   KEY idx_updated_at (updated_at)
  10. ) ENGINE=InnoDB;
  11. -- 目标库(数仓 ODS 层)
  12. CREATE TABLE ods_order (
  13.   id          BIGINT PRIMARY KEY,
  14.   order_no    VARCHAR(32),
  15.   amount      DECIMAL(12,2),
  16.   status      TINYINT,
  17.   updated_at  DATETIME,
  18.   etl_time    DATETIME,
  19.   is_deleted  TINYINT DEFAULT 0
  20. ) ENGINE=InnoDB;
  21. -- 水位表
  22. CREATE TABLE etl_watermark (
  23.   task_name   VARCHAR(64) PRIMARY KEY,
  24.   last_value  DATETIME,
  25.   updated_at  DATETIME
  26. ) ENGINE=InnoDB;
复制代码

步骤 2:方案 A —— 时间戳轮询(带重叠窗口)
  1. -- 1) 取上次水位
  2. SELECT last_value FROM etl_watermark WHERE task_name = 'sync_order';
  3. -- 2) 带 10 分钟重叠窗口抽取(假设上次水位 2026-09-15 02:00:00)
  4. SELECT id, order_no, amount, status, updated_at
  5. FROM t_order
  6. WHERE updated_at >= '2026-09-15 01:50:00'
  7.   AND updated_at <  '2026-09-16 02:00:00';
复制代码

捞出来的数据做幂等 upsert,这一步不能用 insert:
  1. INSERT INTO ods_order (id, order_no, amount, status, updated_at, etl_time)
  2. VALUES (1001, 'NO20260915001', 299.00, 2, '2026-09-15 09:12:33', NOW())
  3. ON DUPLICATE KEY UPDATE
  4.   order_no   = VALUES(order_no),
  5.   amount     = VALUES(amount),
  6.   status     = VALUES(status),
  7.   updated_at = VALUES(updated_at),
  8.   etl_time   = NOW();
复制代码

最后推进水位(注意:水位取的是本次批次的上界,不是 MAX(updated_at)):
  1. INSERT INTO etl_watermark (task_name, last_value, updated_at)
  2. VALUES ('sync_order', '2026-09-16 02:00:00', NOW())
  3. ON DUPLICATE KEY UPDATE last_value = VALUES(last_value), updated_at = NOW();
复制代码

物理删除的兜底做法(如果暂时不上 CDC):每天一次全量主键对账,补标记删除。
  1. -- 目标表有、源表没有的,标记删除
  2. UPDATE ods_order o
  3. LEFT JOIN t_order s ON s.id = o.id
  4. SET o.is_deleted = 1, o.etl_time = NOW()
  5. WHERE s.id IS NULL AND o.is_deleted = 0;
复制代码

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

先确认源库 binlog 格式:
  1. SHOW VARIABLES LIKE 'binlog_format';      -- 必须是 ROW
  2. SHOW VARIABLES LIKE 'binlog_row_image';   -- 建议 FULL
  3. SHOW VARIABLES LIKE 'server_id';          -- 必须非 0 且唯一
复制代码

不对就改配置文件(改完重启实例):
  1. [mysqld]
  2. server_id        = 1001
  3. log_bin          = mysql-bin
  4. binlog_format    = ROW
  5. binlog_row_image = FULL
复制代码

注册连接器(Kafka Connect REST):
  1. {
  2.   "name": "order-connector",
  3.   "config": {
  4.     "connector.class": "io.debezium.connector.mysql.MySqlConnector",
  5.     "database.hostname": "10.0.0.11",
  6.     "database.port": "3306",
  7.     "database.user": "cdc",
  8.     "database.password": "******",
  9.     "database.server.id": "5401",
  10.     "database.include.list": "bizdb",
  11.     "table.include.list": "bizdb.t_order",
  12.     "database.history.kafka.bootstrap.servers": "kafka:9092",
  13.     "database.history.kafka.topic": "schema-changes.bizdb",
  14.     "snapshot.mode": "when_needed"
  15.   }
  16. }
复制代码

消费端按 op 字段(c 新增 / u 更新 / d 删除 / r 快照读)分流处理,d 事件直接打删除标记:
  1. -- 收到 op='d' 的事件时
  2. UPDATE ods_order SET is_deleted = 1, etl_time = NOW()
  3. 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=ROWbinlog_row_image=FULLserver_id 唯一
  • [ ] binlog 保留期 ≥ 最长可接受的中断时长(建议 ≥ 72 小时),否则断连后要重新做快照
  • [ ] CDC 的 server.id 不要和已有从库冲突,否则会导致复制拓扑混乱
  • [ ] DDL 变更单独走审批:加列可以自动兼容,改列名/改类型会导致消费端解析失败
  • [ ] 目标表保留 etl_time 与源端 updated_at 两个时间,问题排查时才能分清是谁的锅
  • [ ] 配置行数波动监控:本次抽取行数偏离近 7 日均值 ±50% 时告警拦停,别让它污染下游


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

—— dbaai
免责申明1、欢迎访问本站,本文内容及相关资源来源于网络,版权归版权方所有!本站原创内容版权归本站所有,请勿转载!
2、本文内容仅代表作者观点,不代表本站立场,作者自负,本站资源仅供学习研究,请勿非法使用,否则后果自负!请下载后24小时内删除!
3、本文内容,包括但不限于源码、文字、图片等,仅供参考。本站不对其安全性,正确性等作出保证。但本站会尽量审核会员发表的内容。
4、如本帖侵犯到任何版权问题,请立即告知本站 ,本站将及时删除并致以最深的歉意!客服邮箱:admin@dbabbs.com
您需要登录后才可以回帖 登录 | 用户注册

本版积分规则

QQ|Archiver|小黑屋|DBA论坛中国 ( 鲁ICP备20017503号-2 )

GMT+8, 2026-9-16 22:16 , Processed in 0.024265 second(s), 10 queries , MemCached On.

Powered by Discuz! X5.0

© 2001-2026 Discuz! Team.

快速回复 返回顶部 返回列表