
Kafka 到下游系统的端到端 Exactly-Once 怎么实现?
拆解"端到端"的三段语义,讲清 Kafka 事务 + 位移原子提交的内环方案、Flink 两阶段提交检查点的外环方案,以及什么时候该退回业务幂等。
共 10 篇文章

拆解"端到端"的三段语义,讲清 Kafka 事务 + 位移原子提交的内环方案、Flink 两阶段提交检查点的外环方案,以及什么时候该退回业务幂等。

从架构模型、吞吐延迟、功能特性(延迟消息/事务/死信)、生态四个维度对比三者,给出日志流选 Kafka、业务消息选 RocketMQ、云原生多租户看 Pulsar 的选型逻辑。

以 lag 监控为入口,按"消费慢、分区不均、Rebalance 频繁、下游瓶颈"四类根因排查,给出临时扩容、加消费者、加分区、跳过历史消息的分级处置方案。

分区数从吞吐目标、消费者并行度、Broker 负载三方面推算,讲清单分区内有序、全局无序的模型,以及按 key 哈希保业务有序的正确姿势。

从磁盘顺序追加、Page Cache 读写缓冲、sendfile 零拷贝、端到端批量压缩四个层面解释 Kafka 的高吞吐,纠正"消息队列都慢在磁盘"的直觉误区。

讲清 Rebalance 的触发条件和三个阶段的流程,说透 Stop The World、重复消费两大影响,以及静态成员、Cooperative 协议等缓解手段。

重复消费的根源是"处理"和"提交位移"不是原子的,拆解自动提交的坑、手动提交的两种模式,给出先处理后提交、业务幂等兜底的工程方案。

幂等解决单分区内的重复写入(PID+序列号去重),事务解决跨分区跨主题的原子写,两阶段提交 + 事务协调器,并能说出各自的边界和代价。

说清 acks 0/1/all 三档语义,给出不丢失的最小配置组合(acks=all + min.insync.replicas + retries),并点破"不丢"是生产、Broker、消费三段的事。

讲清 Broker、Controller、分区副本的角色分工,重点说透 ISR 的入列出列条件、HW 与 LEO 的关系,以及 ISR 为什么能在可靠性和可用性之间取平衡。