我们的交易系统原来用定时任务把 MySQL 数据同步到数据仓库,每 5 分钟跑一次,分析师看到的总是"5 分钟前的旧账"。业务方要实时大屏,等不了。于是用 Kafka 搭了一条 CDC 驱动的实时数据管道,端到端延迟从 5 分钟压到 800 毫秒。
CDC 接入:让数据库自己说变化
定时拉全量太低效,改用 CDC(Change Data Capture):用 Debezium 监听 MySQL 的 Binlog,把每行变更实时发到 Kafka。这样数据管道和源库解耦,不用业务代码改一行:
# Debezium MySQL 连接器(Kafka Connect)
{
"name": "mysql-cdc",
"config": {
"connector.class": "io.debezium.connector.mysql.MySqlConnector",
"database.hostname": "mysql-prod",
"table.include.list": "shop.orders,shop.payments",
"topic.prefix": "cdc"
}
}
每个变更事件带 before/after/op 字段(c/u/d),下游能精确知道"哪行从什么变成什么"。首次全量快照 + 增量 Binlog,平滑衔接,不丢历史。
流处理:用 Kafka Streams 做聚合
原始 CDC 事件是行级变更,业务要的是"每分钟成交额""各渠道占比"。用 Kafka Streams 做有状态聚合:
KStream<String, Order> orders = builder.stream("cdc.shop.orders");
orders.groupBy((k, v) -> v.getChannel())
.windowedBy(TimeWindows.of(Duration.ofMinutes(1)))
.aggregate(() -> 0L,
(k, v, sum) -> sum + v.getAmount(),
Materialized.as("amount-by-channel"))
.toStream()
.to("metrics.channel-amount-1m", Produced.with(...));
状态存在 RocksDB,Kafka 的 changelog topic 负责容错——节点挂了状态能从 changelog 重建。我们实测窗口聚合端到端延迟约 1.2 秒,分析师大屏基本"实时"。
Exactly-Once 如何保障
金融数据最怕重复或丢失。Kafka 从 0.11 起支持幂等生产 + 事务,配合 Streams 的 exactly_once_v2 语义:
props.put(StreamsConfig.PROCESSING_GUARANTEE_CONFIG,
StreamsConfig.EXACTLY_ONCE_V2);
它靠"读处理写"在一个事务里完成:消费 CDC offset、更新 RocksDB 状态、写出聚合结果,要么全成功要么全回滚。我们做过故障注入测试:处理到一半杀掉 Streams 实例,重启后从 checkpoint 恢复,产出结果无重复、无丢失,和手动对账的基准完全一致。
一次对账验证
| 检查项 | 批处理(旧) | Kafka 管道(新) |
|---|---|---|
| 端到端延迟 | 5 分钟 | 0.8 秒 |
| 数据重复率(注入故障后) | 偶发多算 | 0 |
| 源库压力 | 每 5 分钟全表扫 | 仅 Binlog 增量 |
踩坑记录
- Binlog 保留:源库 Binlog 只留 3 天,Debezium 宕机超过 3 天再起就追不上,得全量重做。我们接了监控,Binlog 落后超 1 小时告警;
- Schema 变更:MySQL 加字段,CDC 事件结构变,下游 Streams 反序列化失败。用 Avro + Schema Registry 做兼容演进,新旧 schema 可共存;
- 乱序:网络抖动让 update 事件早于 insert 到达,聚合算错。按事件时间戳(
source.ts)而非到达时间处理,并加 2 秒 watermark 容忍。
小结
Kafka 驱动的实时数据管道,用 CDC 取代定时拉取、用 Streams 做有状态聚合、用事务语义保 Exactly-Once,把"5 分钟旧账"变成"亚秒级实时"。工程要点在三个坑:Binlog 保留与监控、Schema 演进、事件乱序。它不只是技术升级,更是把"数据新鲜度"从分钟级拉到秒级,让业务第一次能基于"现在"而非"刚才"做决策。CDC 接入零侵入这个特性,是它能被快速落地的最大原因。