实时数据链路建设:CDC 采集与流批一体架构

实时数据链路建设:CDC 采集与流批一体架构

业务对数据时效的要求从“次日”变成“小时”再变成“分钟”之后,传统 T+1 批处理链路就捉襟见肘了:一小时的调度窗口里跑不完全量抽取,而下游营销与风控等场景已经在实时看数据。CDC(变更数据捕获)就是为此而生的:不再定时拉全表,而是订阅数据库的变化日志持续增量同步。

数据库变更记录流入流式管道的插画
CDC 的底层逻辑:从“拉取全量”转向“订阅变化”

一、同步方式对比

方式 时效 对源库影响 适用场景
全量抽取(批) 小时至天 高(大表扫描与锁) 小表、维度数据、夜间窗口
增量字段抽取 分钟至小时 中(依赖索引) 有可靠更新时间字段的表
日志型 CDC 秒级 低(读日志) 核心业务表实时同步

日志型 CDC 的优点是几乎不影响业务库,且能捕获 DELETE 与中间状态。代价是强依赖数据库日志能力(如 binlog),运维上需要关注日志保留期与位点管理。

二、四个容易出错的核心问题

1. 一致性:初始快照与增量衔接

CDC 上线时通常要先做一次全量快照,再接增量。难点在于衔接窗口:快照期间发生的变更如果不正确处理,会出现重复或丢失。成熟的实现方式是在快照阶段锁定一致性位点,并将衔接期的变更按幂等方式重放。

2. 幂等与去重

流式链路中“至少一次”是常态,正好一次很难保证。因此下游写入必须幂等:用业务主键 + 版本号(如更新时间戳或 LSN)做条件写入,旧版本不允许覆盖新版本。没有这一层,链路重启一次就会产生脏数据。

3. 顺序性

同一主键的变更必须保序,而并行的好处又想要。可行的做法是按主键哈希分区,保证同一主键的变更进同一个分区;同时下游按版本号比较而非按到达顺序写入,即使乱序到达也不会写错。

4. 大事务与流量突刺

业务上的一次批量操作(如批量导入)在 CDC 侧会变成一瞬间的大量变更,容易把下游压垮。建议在链路中加一层缓冲与背压策略,并对大事务设置阈值告警,必要时改为离线处理。

三、流批一体的现实路径

“流批一体”经常被理解为“一套代码同时跑流和批”。落地时更务实的目标是口径统一与存储统一,代码统一是后续目标:

  • 口径统一:实时指标与离线指标使用同一份指标定义,避免同一个“今日订单数”在两个系统里对不上。
  • 存储统一:湖表或实时数仓作为统一存储,流写入、批回填与修正走同一条路径。
  • 回填能力:必须具备“把历史某段重跑”的能力。流式计算出问题时没有回填手段,是很难接受的。

其中最后一点最容易被忽略,却决定了系统能不能长期运行:设计之初就要想好如何重放数据,否则一次逻辑变更就意味着历史数据永远不一致。

四、监控该盯什么

流式链路的监控与批处理完全不同,重点在“延迟”而不是“成功率”:链路端到端延迟、各算子背压情况、CDC 位点与源库日志的差距(代表是否快跟不上)、数据量黍变与幂等冲突次数。

其中CDC 位点滞后量是最关键的健康指标:它反映消费速度是否持续跑得上生产速度。一旦这个数字单调上漲,即使当前查询看似正常,也已经是故障前兆。

五、落地建议

不要一次性把所有表都切到 CDC。建议先选一张写入量中等的核心表(如订单主表),完成从采集、清洗、入仓到指标输出的全链路,把重启、回填、去重这些异常场景验证完,再按表批量扩展。实时链路的复杂度主要来自异常路径而非正常路径,先把它跑熟比多接十张表更有价值。

Related

相关阅读

这篇文章讨论的问题,我们也许能帮你解决

把你的场景描述给我们,一起看看有没有更省的路径。