Administrator
发布于 2024-07-10 / 959 阅读
19

Kafka 在实时数据管道中的架构实践

我们的交易系统原来用定时任务把 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 接入零侵入这个特性,是它能被快速落地的最大原因。

参考