千万级数据同步与 ETL 实战(六):CDC、事件时间与真正的流式处理
前五篇使用高频 Micro Batch 处理权限受限的数据源。它能做到分钟级同步和故障重放,但事件仍由调度器分批查询。真正的流式链路持续消费无界事件,数据一旦进入日志或消息通道便可推动计算,系统还要处理分区顺序、事件时间、乱序、长时间状态、Barrier 快照和下游事务。
这套架构适合秒级或亚秒级延迟、持续高变化率、窗口聚合、复杂乱序以及需要按日志位置精确恢复的场景。源端只有分页查询、文件投递或低频 API 时,部署 Flink 并不会凭空生成事件流,输入端仍受轮询边界约束。
一、流式输入需要可持续的位置
数据库 CDC 读取提交日志。MySQL 以 Binlog 文件与位置或 GTID 标识进度,PostgreSQL 通过逻辑复制槽和 LSN 暴露变化,Oracle 常使用 SCN 关联日志顺序。连接器将底层记录转换为统一事件,同时保留源端位置:
{
"source": "orders",
"partition": "mysql-cluster-a",
"position": {"gtid": "synthetic:1-84321"},
"transactionId": "tx-1057",
"operation": "UPDATE",
"key": {"orderId": 10086},
"before": {"status": "PAID"},
"after": {"status": "SHIPPED"},
"sourceTime": "2026-07-19T10:01:03.126Z",
"schemaId": "orders-v9"
}日志位置承担恢复游标,事务 ID 和事务边界承担一致性语义。数据库事务修改 100 行时,连接器需要保留它们的相对顺序,并明确下游是否按整笔事务可见。将每行直接发往不同分区,吞吐会提高,却可能让下游在短时间内观察到部分事务结果。
日志保留时间必须覆盖停机和恢复窗口。消费任务离线 12 小时,而源端日志只保留 6 小时,旧位置已不可用,只能重新做快照并建立新起点。复制槽也不是免费保险:消费者停滞时,PostgreSQL 可能长期保留 WAL,造成磁盘增长。监控应同时展示消费延迟和源端日志保留压力。
无法开放数据库日志时,可以由业务事务写 Outbox 表。业务数据与 Outbox 事件在同一数据库事务内提交,CDC 只读取 Outbox。跨数据库、外部 HTTP 调用和消息发送无法靠一张表自动形成原子操作,仍需幂等、补偿或业务 Saga。
二、顺序只在定义的范围内成立
单个 Kafka Partition 内有稳定偏移顺序,多个 Partition 之间没有全局顺序。按业务键分区可保证同一订单的事件进入同一 Partition:
$$ partition=hash(orderId)\bmod P $$
处理实例扩缩容时,Partition 会重新分配,但每个 Partition 的 Offset 仍按序消费。若同一业务实体使用过不同分区键,或上游在重分区过程中改变了 Hash 规则,同键顺序就可能被打散。
全局排序通常要把所有事件集中到一个分区,吞吐上限和故障域都会集中。多数业务只需要键内顺序。跨键约束应转换成有状态协调、事务边界或后续校验,不要在需求文档中写一个没有范围的“消息有序”。
CDC 还会受到表级和事务级并行影响。单库日志可以提供提交序列,分库分表后每个库有独立位置。若业务事务不跨库,按实体键合并即可;若要求跨库全局一致视图,需要额外的全局版本、协调器或业务层事件协议。
三、事件时间和处理时间回答不同问题
流任务里至少有两个时间:
event time 事件在业务或源系统中发生的时间
processing time 算子处理事件时的机器时间处理时间窗口实现简单,但结果受网络延迟、重启和背压影响。同一批历史事件重放时,处理时间已经变化,窗口结果难以复现。事件时间窗口按记录携带的时间归属,适合订单量、设备指标和日志统计。
事件时间也可能不可靠。上游时钟偏差、客户端离线、字段精度和业务回填都会造成乱序。系统需要 Watermark 表示“预计不会再看到早于某时刻的大量事件”。对于分区 $i$,可以根据已观察到的较大事件时间 $T_i^{max}$ 和乱序容忍 $\delta_i$ 生成:
$$ W_i=T_i^{max}-\delta_i $$
多输入算子的 Watermark 通常受慢分区限制:
$$ W_{operator}=\min_i W_i $$
如果某个分区长时间没有数据,它会阻塞整个算子的事件时间推进。运行时需要 Idle Partition 检测,将确认空闲的分区暂时移出较小值计算。空闲超时过短会把随后到达的旧事件判为迟到,过长会让窗口长时间不关闭,应根据分区流量间隔分布设置。
Watermark 代表启发式进度,不是墙上时钟,也不是“该时间前绝无事件”。迟到事件必须有处理策略。
四、窗口由归属、触发和清理共同定义
以 5 分钟滚动窗口为例,事件时间 $t$ 的窗口起点为:
$$ start=\left\lfloor\frac{t}{5m}\right\rfloor\times5m $$
窗口范围使用左闭右开区间 [start, start+5m)。当 Watermark 越过窗口结束时间,算子触发计算。若允许迟到 2 分钟,状态至少保留到:
$$ cleanupTime=windowEnd+2m $$
迟到但未超过保留期的事件可以更新结果,下游要支持 Update 或 Retract。超过保留期的事件进入 Side Output、补偿流或审计队列,不能静默丢弃。
滑动窗口由窗口大小 $S$ 和滑动步长 $s$ 决定,一条事件可能属于 $\lceil S/s\rceil$ 个窗口。S=1h、s=1m 时,每条事件会进入约 60 个逻辑窗口。直接复制状态会增加写放大,运行时可使用 Pane 预聚合,再组合为滑动窗口结果。
Session Window 按键维护活动会话,事件间隔超过 Gap 才关闭。乱序事件可能连接两个已存在会话,需要支持窗口合并。会话状态增长与活跃 Key 数、乱序上限和 Gap 共同相关,不能只按每秒事件数估算。
五、Keyed State 决定内存和恢复规模
有状态算子把数据按 Key 分区,每个 Key 保存聚合、去重、定时器或模式匹配状态。状态量可粗略估算为:
$$ StateBytes\approx ActiveKeys\times BytesPerKey\times VersionFactor $$
1000 万活跃 Key、每 Key 200 字节裸状态就是约 2GB;加入序列化、索引、快照版本和本地存储开销后会更高。窗口重叠、去重集合和长 TTL 还会继续扩大状态。
State TTL 负责回收长期不再访问的 Key,但 TTL 到期语义需要和业务一致。订单状态保存 7 天后被清理,8 天前的迟到更新会被当作新状态处理。高基数去重若要求永久精确,状态会持续增长,可以改为按业务保留期分桶,或接受 Bloom Filter 等概率结构的误判边界。
状态分区要与输入 Key 一致。修改分区数时,Savepoint 或 Checkpoint 中的 Key Group 会重新分配给新并行实例。自定义状态若绕过运行时管理,扩缩容后很难重新分片和恢复。
六、Checkpoint Barrier 如何得到一致状态
分布式流图中的算子并行处理多条输入。Checkpoint Coordinator 向 Source 注入编号为 $c$ 的 Barrier。Barrier 随数据流动,区分 Checkpoint 前后的记录。算子收到所有输入的 $c$ Barrier 后,保存状态并向下游转发。
多输入算子的对齐 Checkpoint 会暂停已经收到 Barrier 的通道,等待其余通道到达。这样快照不需要保存 Barrier 之后的在途记录,但慢通道和背压会增加对齐时间。非对齐 Checkpoint 将在途 Buffer 一并保存,减少对齐等待,代价是 Checkpoint 体积和恢复读取量上升。
Checkpoint 周期越短,故障后重放越少,正常运行的持久化开销越高。设事件处理速率为 $r$,Checkpoint 间隔为 $I$,故障时均匀落在区间内,预期重放事件数量约为:
$$ E[replay]\approx\frac{rI}{2} $$
还要计入 Checkpoint 完成耗时。若一次快照长期无法在下一次触发前完成,应先处理状态体积、存储带宽或背压问题,不能只继续缩短间隔。
七、端到端 Exactly Once 的成立条件
流处理框架保存状态一次,不等于外部系统只收到一次。端到端一致性要求把三个位置绑定到同一个 Checkpoint:
Source Position
+ Operator State
+ Sink CommitSource 恢复到 Checkpoint 记录的 Offset,算子恢复对应状态,Sink 只提交该 Checkpoint 对应的事务。Sink 支持两阶段提交时,可以在 Barrier 到达时预提交,Checkpoint 全局完成后提交事务;任务失败则中止未完成事务。
外部 HTTP API 通常没有两阶段事务。此时使用稳定幂等键,例如 (jobId, checkpointId, partition, sequence),服务端保存处理结果并在重复请求时返回原结果。若 API 既无事务也无幂等接口,框架无法提供严格端到端一次效果,只能通过重试、查询确认和对账降低重复风险。
Kafka 事务可以原子提交输出记录和已消费 Offset,但保证范围限于参与同一事务协议的 Kafka 资源。写数据库、对象存储和第三方接口时,需要分别核对连接器语义。文档里应写“Kafka 到 Kafka 事务链路”或“数据库 Sink 幂等 Upsert”,不要使用无范围的 Exactly Once 标签。
八、背压用排队关系解释
当持续输入速率 $\lambda$ 大于处理速率 $\mu$ 时,队列增长率近似为:
$$ \frac{dQ}{dt}=\lambda-\mu,\quad \lambda>\mu $$
队列增长会延长端到端延迟,并让 Checkpoint Barrier 更难推进。稳定系统的长期平均必须满足 $\lambda<\mu$,还要保留处理突发流量和恢复重放的余量。
Little 定律给出稳定区间内的关系:
$$ L=\lambda W $$
$L$ 表示系统内平均在途事件数,$W$ 表示平均停留时间。输入保持每秒 5 万条且平均停留 20 秒时,在途规模约为 100 万条。只看算子 CPU 而不看 Busy Time、反压时间、Buffer 水位和 Checkpoint 对齐耗时,很难判断延迟增加发生在哪一段。
反压需要传播到 Source。Kafka Source 可暂停拉取,但 Offset Lag 会增长;数据库 CDC Source 不能无限阻止数据库产生日志,停滞过久会增加日志保留压力。容量规划应覆盖峰值输入、稳定处理能力、故障恢复速率和允许积压时长。
热点 Key 会让单个分区先进入反压。随机加盐可以拆散热点聚合,但随后需要二阶段聚合,且严格键内顺序会被改变。状态更新有顺序依赖时,应先优化单 Key 算法、拆分业务实体或隔离热点,不能只改分区键。
九、Schema 演进要支持历史重放
CDC 事件携带 Schema 版本。新增可空字段通常可向前兼容。重命名会破坏字段映射,数值精度缩小可能拒绝旧值,主键调整还会改变分区和状态归属,这些变更都需要专门的迁移版本。流任务恢复或重放旧日志时,仍会读到历史 Schema,反序列化器和转换逻辑必须保留相应版本。
Schema Registry 需要定义兼容策略:
Backward:新消费者能读取旧数据
Forward:旧消费者能读取新数据
Full:同时满足两者
None:每次变更由发布流程显式处理数据库 DDL 事件和数据事件的顺序也要保留。连接器发现新列后,下游 Sink 表可能尚未扩展。可先发布兼容目标 Schema,再允许源端 DDL;无法协调时,事件进入带 Schema 的缓冲区,完成迁移后重放。
Savepoint 升级还涉及状态序列化兼容。算子 UID 保持稳定,状态描述符和序列化器变更要有迁移器。直接改类名、字段类型或算子拓扑,可能让旧状态无法恢复。升级演练应从生产规模的脱敏 Savepoint 启动作业,并执行一段历史事件回放。
十、真流式和微批可以共存
CDC 通道提供低延迟,但仍可能因日志过期、连接器缺陷、错误 Schema 或人工绕过业务写入而产生差异。前五篇的 Warm 和 Cold 通道可以作为流式系统的校验面:
三条通道共享身份键、Schema 和版本协议,Checkpoint 各自独立。Warm 通道发现流式链路遗漏后,生成带扫描证据的修复事件;目标端按源版本防止旧扫描结果覆盖新 CDC 事件。Cold 通道低频重建摘要,验证两侧状态是否收敛。
选择方案时可以用延迟、输入能力和计算语义判断:
| 条件 | 合适的主链路 |
|---|---|
| 只读查询、分钟级延迟 | Watermark Micro Batch |
| 无更新时间、允许小时级发现 | Hash Wheel + Fingerprint |
| 可读取日志、只做状态同步 | CDC + 幂等 Sink |
| 有事件时间窗口和乱序计算 | 有状态流处理 |
| 文件或 API 定期到达 | 版本化 Micro Batch |
| 强实时主链路还需长期对账 | CDC Hot + 扫描 Warm/Cold |
十一、测试运行时语义
测试数据要主动制造乱序、重复、迟到、空闲分区、热点 Key 和 Schema 切换。固定输入与随机重启点,分别在 Barrier 前、状态快照后、Sink 预提交后和全局完成通知前终止进程,检查恢复结果。
事件时间测试由测试用例显式推进虚拟时钟和 Watermark,机器墙钟不参与窗口触发。窗口结果要验证首次输出、迟到更新、清理后 Side Output 和重放一致性。Exactly Once 测试应查询目标端事务或幂等登记,确认重复尝试存在但业务结果只生效一次。
容量测试持续运行到状态、Checkpoint 和日志保留进入稳态。短时间压测看不到 TTL 清理、Compaction、增量 Checkpoint 退化和长尾分区。故障恢复期间还要保持新流量输入,测量追平速度:恢复处理能力若只等于当前输入速率,积压永远无法清空。
流式系统的难点集中在可恢复的时间和状态。日志位置定义从哪里继续,Watermark 决定何时关闭事件时间窗口,Checkpoint Barrier 固化分布式状态,Sink 协议决定故障后重复是否影响业务。把这些边界逐项验证后,秒级链路才具有可解释的正确性。