监控显示 Consumer Lag 为 0,客服却仍看到半小时前的库存。Lag 没有撒谎:它只说明 Group 的已提交 offset 追上 Log End Offset,不证明数据库写入、索引刷新和缓存失效都完成。
Kafka Lag 是传输位置指标,不是端到端业务新鲜度 SLA。
一条记录至少经过五个位点
Kafka LEO → Consumer fetched position → committed offset → business processing completed → downstream visible versionkafka-consumer-groups.sh --describe展示CURRENT-OFFSET、LOG-END-OFFSET与二者差值。官方将它定义为 Consumer 位置检查工具;输出没有“数据库已提交”或“用户已看到”字段。Basic Operations
四种 Lag=0 的异常路径
| 路径 | Kafka 侧 | 业务侧 |
|---|---|---|
| 先提交 offset,后异步处理 | Lag 很快归零 | 队列仍积压或任务失败 |
| 异常被吞掉仍提交 | Lag 归零 | 部分记录永久未落地 |
| 数据库已更新,索引/缓存滞后 | Lag 归零 | 查询仍返回旧版本 |
| 消费了错误 Topic/环境 | 目标 Group 正常 | 用户链路根本未被更新 |
相反,Lag 大也不一定等于用户数据旧:若积压的是低优先级 Key,而当前热点数据已被更新,业务新鲜度可能仍达标。
Commit 不是业务完成证明
Consumer 维护当前位置,也可以把 offset 提交给 Group Coordinator 以便重启恢复。Distribution 自动或手工提交只是保存恢复起点;如果提交发生在异步任务完成前,Group 看起来健康,失败任务却可能在重启后被越过。
可靠处理至少要明确:
- offset 何时提交;
- 批次内部分失败是否阻止越过;
- 外部写入是否幂等;
- 重试队列、死信或补偿是否纳入新鲜度统计。
只读排查:从用户看到的业务键反向追踪
1. 记录一个具体样本
不要从总 Lag 猜原因。选取一个用户可复现的业务键,收集event_id、业务版本、事件时间、Topic、partition、offset、目标库版本与缓存版本。
2. 核对 Group 位点
bin/kafka-consumer-groups.sh --bootstrap-server broker:9092\--describe--groupinventory-projection只读观察。正常仅代表 committed offset 接近 LEO;若目标记录 offset 小于 CURRENT-OFFSET 而下游没有对应版本,证明“位点已越过、结果未闭环”。
3. 对齐应用阶段耗时
应用必须分别记录:poll 时间、处理开始、外部提交、offset commit、缓存失效完成。只记录“消费成功”无法区分哪个阶段成功。
4. 核对最终读取路径
直接读事实库、搜索索引和缓存,比较业务版本而非机器时间。若事实库新而缓存旧,修 Kafka 没有意义;若事实库也旧,再查消费异常和提交时序。
Kafka 官方建议监控records-lag-max与最小 fetch rate,但端到端系统还必须补充处理队列深度、外部提交失败率和业务版本年龄。Monitoring
处置与止损
- 暂停继续越过失败记录的 commit 路径,保留失败样本和原始 offset。
- 修复异步完成确认:只有批次中连续完成的 offset 才能推进提交。
- 目标端按事件 ID/业务版本幂等,允许安全重放。
- 缓存和索引建立独立积压与新鲜度指标,不再借用 Kafka Lag。
- 对已漏记录生成精确重放清单,避免整组无边界回退。
Offset reset 会改变消费历史。执行前必须停止 Group 活动实例、保存各分区原 offset、先预览目标、限定 Topic/分区和时间窗;出现重复副作用或目标库压力超阈值立即停止,并可恢复原 offset。
业务新鲜度应该怎样度量
更有用的是端到端年龄:
freshness_age = now - latest_business_version_visible_time同时按阶段拆分:Kafka 等待、应用等待、外部提交、索引刷新、缓存传播。技术验收要求每个事件的状态链闭合;业务验收要求用户查询在 SLA 内返回不低于事件业务版本的结果。
源码与 Java:只提交真正完成的位点
以下源码定位与 Java 示例按 Kafka 4.3.1 静态审阅,未在本环境运行;发布前请在隔离 Topic 和测试 Group 中验证依赖、权限与业务幂等语义。
KafkaConsumer.poll/commitSync分别推进本地 position 与 Group committed offset;服务端提交落到OffsetMetadataManager。源码里没有“数据库已可见”状态。
importjava.time.Duration;importjava.util.*;importorg.apache.kafka.clients.consumer.*;importorg.apache.kafka.common.TopicPartition;importorg.apache.kafka.common.serialization.StringDeserializer;publicclassCommitAfterBusinessSuccess{staticvoidwriteBusiness(Stringkey,Stringvalue){/* 用业务唯一键执行幂等数据库事务 */}publicstaticvoidmain(String[]args){Propertiesp=newProperties();p.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG,"localhost:9092");p.put(ConsumerConfig.GROUP_ID_CONFIG,"inventory-projection");p.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG,"false");p.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG,StringDeserializer.class);p.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG,StringDeserializer.class);try(KafkaConsumer<String,String>c=newKafkaConsumer<>(p)){c.subscribe(List.of("inventory-events"));ConsumerRecords<String,String>rs=c.poll(Duration.ofSeconds(10));Map<TopicPartition,OffsetAndMetadata>done=newHashMap<>();for(ConsumerRecord<String,String>r:rs){writeBusiness(r.key(),r.value());done.put(newTopicPartition(r.topic(),r.partition()),newOffsetAndMetadata(r.offset()+1));}c.commitSync(done);}}}映射是poll → position、业务成功后commitSync → OffsetMetadataManager → CURRENT-OFFSET。示例要求writeBusiness自身幂等;它仍不能证明缓存或索引已刷新,因此必须用业务版本继续验收。
结论
Lag=0 只证明位点追平,不证明业务完成。真正的生产闭环要用事件 ID 和业务版本贯穿 Kafka、应用、数据库、索引与缓存,让“哪里旧、旧多久、能否重放”都可回答。