三亩地 三亩地SAN MU DI · CODE DIARY
ARTICLE DETAIL

日记详情

真实记录编程学习的某一天,欢迎挑你感兴趣的翻一翻。

Kafka只管流不管应答——内存库异步落库的边界与避坑

Kafka只管流不管应答——内存库异步落库的边界与避坑

Kafka只管流不管应答——内存库异步落库的边界与避坑

Kafka 擅长把数据从 A 搬到 B,顺带让 C、D 也各搬一份;但它不擅长"你问我答"。把它当同步 RPC 用,是这类中间件最常见的误用。这篇不讲 Kafka 怎么安装,讲一个政务老兵用它做"内存库 → MySQL 异步落库"时踩过的边界和坑。

文章目录

  • Kafka只管流不管应答——内存库异步落库的边界与避坑
    • 一、先认清 Kafka 是什么
    • 二、消息会过期:最大的隐形风险
    • 三、生产者:简单到让人不安
    • 四、主线场景:内存库 → Kafka → MySQL
      • 这个架构值在哪
      • 四条硬性开发要求
    • 五、踩坑:消费组停机一周,重启后数据缺了一块
    • 六、横向对比:Kafka vs MSMQ
    • 七、避坑清单(六条)
    • 八、一句话浓缩

一、先认清 Kafka 是什么

很多人上手 Kafka 第一件事就是去找"怎么拿到回执"——找不到,于是觉得它难用。这不是 Kafka 难用,是用错了地方。

Kafka 是异步事件流中间件,天生单向流转。生产者把消息丢进 Topic 就返回了,不等任何人消费。它不适合同步响应式请求(一问一答的 RPC 场景)。如果你的业务是"我发一条,必须马上拿到下游处理结果才能往下走"——别用 Kafka,用 HTTP/ RPC。

认清这一条,后面的选型就不会跑偏。

核心机制靠三样东西撑起来:Topic、分区 Partition、消费组 Group

  • 同一个 Group 下挂多个消费者:负载均衡,一条消息只会被一个消费者处理。并发的上限 = 分区数,加再多消费者也白搭。
  • 不同 Group 各自订阅同一个 Topic:效果近似广播,各组独立消费全量数据。
  • Broker不记录"哪条消息被读过",只在内置 Topic__consumer_offsets里持久保存【消费组 + Topic + 分区】对应的 offset。offset 由消费者自己主动提交。

最后这一条是 Kafka 和传统 MQ 最大的分水岭:消息只存一份,谁消费到哪了自己记。这个设计让它的存储极友好,但也埋了坑——下面第二节讲。


二、消息会过期:最大的隐形风险

这是新手最容易忽略、也最容易翻车的一条。

Topic 内的消息有默认过期策略(7 天,log.retention.hours=168,达到时间阈值或容量阈值,Kafka 会直接删掉旧的日志段。

关键规则:

  • 消息一旦被清理,所有消费组都无法再读取——不管你有没有消费过。
  • Kafka 只保存一份消息日志,多消费组只是各自维护消费位置,不会复制消息。

👉 风险场景:消费服务长时间停机(比如国庆长假停机维护 8 天),超过消息留存时间,这段时间的增量数据永久丢失,重置 offset 也找不回来。

这条风险在政务系统里特别要命——社保缴费、医保结算的数据丢一天都是事故。所以做异步落库时,消息留存时间一定要根据"最长可能停机时间"来调,默认 7 天在很多政务场景里根本不够。


三、生产者:简单到让人不安

Kafka 生产者这一侧简单得有点过分:

  • 完全无感知下游:生产者不知道有几个消费者、几个消费组,发消息时不需要任何特殊改造。
  • 想保证相同业务数据有序:发送时指定业务 Key(比如身份证号、订单号),消息会按 Key 哈希固定进入同一个分区,同 Key 永远有序。
  • 可靠性配置acks、幂等、重试)属于业务需求层面的开关,和下游消费架构无关。

这种"上游不管下游"的设计,本质上是解耦的极致——但也意味着出了问题,上游不会主动通知你。所以消费端的健壮性必须自己兜。


四、主线场景:内存库 → Kafka → MySQL

这是我用 Kafka 最顺手的场景,也是它真正发光的地方。

背景:政务高并发业务(比如集中参保期缴费),业务层先把数据写进内存数据库(我自己的 CacheSQL,或 Redis、或内存表),保证毫秒级响应;然后通过 Kafka 异步把数据落到MySQL/Oracle持久化。

业务请求 │ ▼ 内存数据库(毫秒级读写,扛峰值) │ ▼ 写入 Topic(指定身份证号做 Key) Kafka(削峰 + 解耦 + 可回溯) │ ├─ 消费组A → MySQL(持久化主链路) ├─ 消费组B → ElasticSearch(建搜索索引) └─ 消费组C → 数据仓库(报表分析)

这个架构值在哪

  1. 削峰:峰值每秒几万次内存写入,Kafka 顺序落盘扛得住,消费端按自己的节奏慢慢入库。
  2. 上下游解耦:内存库只管写 Kafka,不关心几个下游、各自怎么落库。
  3. 故障回溯:消费端出问题,重置 offset 回到故障点重新消费——这是 Kafka 区别于传统 MQ 的杀手锏。
  4. 一份事件多下游消费:多消费组各取所需,内存库里只产生一份数据。

四条硬性开发要求

这个架构能用,但有四条铁律不能破:

① 手动提交 offset,禁止自动提交。

// 反例:自动提交,消费方法还没跑完 offset 就提交了// 一旦消费失败,这条消息等于没处理,直接丢props.put("enable.auto.commit","true");// 正解:关掉自动提交,处理完业务逻辑再手动提交props.put("enable.auto.commit","false");consumer.commitSync();// 只在业务落库成功后才调用

自动提交是 Kafka 默认行为,也是最容易丢数据的配置——它按时间间隔提交,不管你这条消息处理完没。在涉及基金数据的政务场景,这条必须关。

② 必须实现消费幂等。

Kafka 默认是 at-least-once(至少一次)语义,意味着同一条消息可能被消费不止一次。落库时必须按业务主键做幂等:

-- 用 INSERT ... ON DUPLICATE KEY UPDATE 兜重复-- 或先 SELECT 判断存在再 INSERT-- 总之:同一条消息消费 10 次,数据库里也只有 1 条

③ 相同主键的消息指定 Key,规避乱序。

同一个人的多条变更(比如先改姓名、再改缴费档次),必须按身份证号哈希到同一分区,保证顺序。不指定 Key,消息会被轮询打散到不同分区,先改的反而后落库。

④ 接受异步延迟,别想毫秒级强一致。

内存库和 MySQL 之间有几秒到几十秒的延迟是正常的——这就是异步的代价。如果业务要求"写入立即能在 MySQL 查到",那不该走 Kafka,应该同步双写。认清边界,该同步就同步,别难为 Kafka。


五、踩坑:消费组停机一周,重启后数据缺了一块

这是真实出过的事,比"正常流程"更值得记。

现象:参保缴费的 Kafka 消费服务因为国庆长假停机维护,节后重启,发现 10 月 1 日到 10 月 3 日这三天的缴费记录在 MySQL 里查不到,参保人反映"明明扣了款,系统里查不到缴费"。

定位:先查消费日志,offset 从 10 月 4 日开始正常推进,没有报错。再看 Kafka Topic 的消息留存配置——log.retention.hours=168(7 天)。问题是长假停了 8 天,10 月 1-3 日的消息日志段已经被 Kafka 按过期策略清理掉了。

原因:消费组长时间不运行,超过了 Topic 的消息留存时间。Kafka 不会因为"还有人没消费"就保留消息——过期就删,这是它的存储策略,跟有没有人消费无关。offset 重置到最早位置也没用,消息物理上已经不存在了。

方案

  • 应急:从内存库的事务日志里把这三天的事务重放一遍(幸好内存库有 WAL),补进 Kafka 重新消费。这次没造成基金损失,但吓出一身冷汗。
  • 根治
    • 把 Topic 留存时间从 7 天调到30 天log.retention.hours=720),覆盖最长可能的停机窗口。
    • 消费服务配置死信队列,消费失败的消息转存而不是丢。
    • 关键业务(涉及基金的)增加一道对账兜底:每天比对内存库当日事务数和 MySQL 落库数,不一致告警。

验证:调整留存时间后,故意停掉消费服务 10 天再启动,offset 回溯正常,消息完整可消费。

边界:Kafka 的消息留存不是无限期的,它是"存储友好"和"可靠性"之间的权衡。你的留存窗口必须 ≥ 你的最长可能停机时间,否则就是赌命。


六、横向对比:Kafka vs MSMQ

政务系统里经常能碰到 MSMQ(Windows 原生消息队列),它能做"内存库同步数据库"这事的简化版,但能力边界差很多:

维度KafkaMSMQ
跨平台跨平台,Linux 主流绑定 Windows
吞吐高(顺序落盘,每秒数万~数十万)一般
消息模型一份消息,多消费组各自维护 offset取出即移除,原生不支持多业务独立消费同一份
回溯能力支持,重置 offset 即可重放没有,取出就没了
横向扩展分区 + 消费组,扩展能力强

选型建议

  • 小规模 Windows 传统系统、单一下游、不需要回溯 → MSMQ 够用,别为它单独搭 Kafka 集群。
  • 大数据量、多个下游各自消费、需要故障回溯 → 优先 Kafka。

七、避坑清单(六条)

  1. 单纯增加消费者不能提升并发——分区数决定上限。4 个分区的 Topic,挂 8 个消费者,有 4 个永远闲置。
  2. 消费组长时间不运行有两个风险——消息过期删除、offset 本身被系统清理(offsets.retention.minutes默认 7 天)。
  3. 自动提交 offset 极易引发数据丢失——业务没处理完就提交了,必须关掉。
  4. 不要用 Kafka 搭同步请求-应答架构——可以靠临时 Topic + correlationId 凑出来,但架构复杂、缺陷多,得不偿失。
  5. “多消费组订阅同一 Topic” ≠ 传统 MQ 的广播——受消息生命周期限制,不能无限存放历史数据,晚加入的消费组拿不到加入前的消息。
  6. 留存时间必须覆盖最长停机窗口——这条单独拎出来,因为它最致命。

八、一句话浓缩

Kafka 是面向数据流的异步消息引擎,擅长事件分发和数据异步同步;靠消费组实现负载均衡与多订阅广播,支持故障回溯;但存在消息过期限制、不适合同步调用。选型时要区分 Windows 原生 MSMQ 的能力边界,开发层面必须做好手动提交 offset 与幂等处理。

认清边界,用对的场景,它就是解药;用错场景,它就是数据黑洞。

← 返回列表