你们项目为什么用 Kafka?一条核心业务消息从发送到消费成功,如何保证端到端可靠?
这个问题我会先说项目结论:串联生产、Broker、消费与数据库,设计可恢复的端到端投递链路
我会先交代项目背景和选型结论。消息系统常见语义包括最多一次、至少一次和恰好一次。工程上最常见的是“至少一次投递 + 业务幂等”,因为跨 Kafka、数据库和外部服务的端到端恰好一次需要所有环节共同参与。
具体落地时,我会沿着实际调用链来讲。沿着「生产者发送并等待确认 → Leader 追加日志 → Follower 拉取并进入 ISR → 高水位推进可见性 → 消费者处理并提交位点」观察输入、状态与输出,这些阶段都可以从日志、指标或源码里验证。 生产者发送并等待确认 acks、重试和幂等生产者决定客户端如何确认写入;超时可能代表失败也可能代表结果未知。 Leader 追加日志 Leader 先把记录追加到本地分区日志,顺序仅在单分区内定义。 ReplicaManager:Leader 写与 ISR。 kafka-topics.sh --describe:副本/ISR 分布。 我会先确认请求实际走到了哪条路径,再用运行数据验证,不会只看类名或配置猜测。
参数和容量不能靠默认值,我会结合业务量来定。acks=all 要求 Leader 等待同步副本确认,可靠性高于只等待 Leader。它需要与合理的副本数、min.insync.replicas 配合,否则副本不足时仍可能降低容灾能力或导致写入失败。 生产者应开启幂等能力,使用有限重试并监控最终失败。幂等生产者可以避免重试在单分区内产生重复写入,但不能自动让下游数据库操作变成幂等。
效果要用数据证明,线上问题也要能闭环。固定输入和基线 先在没有故障注入的环境执行上述配置,固定数据规模、并发度、运行时版本和预热时间。以「UnderReplicatedPartitions」为主基线,记录值应满足「记录分区基线」;同时保存 UnderReplicatedPartitions、ISR 收缩,使后续变化能够回到同一时间轴比较。 Topic 有三副本,但生产者使用 acks=1,Leader 本地确认后尚未复制就宕机。副本数不等于写确认强度;改为 acks=all、min.insync.replicas=2 并禁用不干净选主后才符合 RPO。
最后我会主动说明这个方案不适合什么场景。方案:更适合的场景:主要收益:代价与边界。 acks=1:允许极小丢失窗口且延迟敏感:确认快:Leader 故障可能丢已确认写。 acks=all + min ISR:核心事件:副本确认更可靠:副本不足时拒绝写入。 端到端 Outbox + 幂等:数据库状态与消息需一致:跨系统有可恢复证据:链路和运维更复杂。 选型至少带上 消息速率、峰值带宽、分区数、消息大小和积压恢复时间,并用上面的量化基线验证;