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.hours168达到时间阈值或容量阈值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 → 数据仓库报表分析这个架构值在哪削峰峰值每秒几万次内存写入Kafka 顺序落盘扛得住消费端按自己的节奏慢慢入库。上下游解耦内存库只管写 Kafka不关心几个下游、各自怎么落库。故障回溯消费端出问题重置 offset 回到故障点重新消费——这是 Kafka 区别于传统 MQ 的杀手锏。一份事件多下游消费多消费组各取所需内存库里只产生一份数据。四条硬性开发要求这个架构能用但有四条铁律不能破① 手动提交 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.hours1687 天。问题是长假停了 8 天10 月 1-3 日的消息日志段已经被 Kafka 按过期策略清理掉了。原因消费组长时间不运行超过了 Topic 的消息留存时间。Kafka 不会因为还有人没消费就保留消息——过期就删这是它的存储策略跟有没有人消费无关。offset 重置到最早位置也没用消息物理上已经不存在了。方案应急从内存库的事务日志里把这三天的事务重放一遍幸好内存库有 WAL补进 Kafka 重新消费。这次没造成基金损失但吓出一身冷汗。根治把 Topic 留存时间从 7 天调到30 天log.retention.hours720覆盖最长可能的停机窗口。消费服务配置死信队列消费失败的消息转存而不是丢。关键业务涉及基金的增加一道对账兜底每天比对内存库当日事务数和 MySQL 落库数不一致告警。验证调整留存时间后故意停掉消费服务 10 天再启动offset 回溯正常消息完整可消费。边界Kafka 的消息留存不是无限期的它是存储友好和可靠性之间的权衡。你的留存窗口必须 ≥ 你的最长可能停机时间否则就是赌命。六、横向对比Kafka vs MSMQ政务系统里经常能碰到 MSMQWindows 原生消息队列它能做内存库同步数据库这事的简化版但能力边界差很多维度KafkaMSMQ跨平台跨平台Linux 主流绑定 Windows吞吐高顺序落盘每秒数万~数十万一般消息模型一份消息多消费组各自维护 offset取出即移除原生不支持多业务独立消费同一份回溯能力支持重置 offset 即可重放没有取出就没了横向扩展分区 消费组扩展能力强弱选型建议小规模 Windows 传统系统、单一下游、不需要回溯 → MSMQ 够用别为它单独搭 Kafka 集群。大数据量、多个下游各自消费、需要故障回溯 → 优先 Kafka。七、避坑清单六条单纯增加消费者不能提升并发——分区数决定上限。4 个分区的 Topic挂 8 个消费者有 4 个永远闲置。消费组长时间不运行有两个风险——消息过期删除、offset 本身被系统清理offsets.retention.minutes默认 7 天。自动提交 offset 极易引发数据丢失——业务没处理完就提交了必须关掉。不要用 Kafka 搭同步请求-应答架构——可以靠临时 Topic correlationId 凑出来但架构复杂、缺陷多得不偿失。“多消费组订阅同一 Topic” ≠ 传统 MQ 的广播——受消息生命周期限制不能无限存放历史数据晚加入的消费组拿不到加入前的消息。留存时间必须覆盖最长停机窗口——这条单独拎出来因为它最致命。八、一句话浓缩Kafka 是面向数据流的异步消息引擎擅长事件分发和数据异步同步靠消费组实现负载均衡与多订阅广播支持故障回溯但存在消息过期限制、不适合同步调用。选型时要区分 Windows 原生 MSMQ 的能力边界开发层面必须做好手动提交 offset 与幂等处理。认清边界用对的场景它就是解药用错场景它就是数据黑洞。