跳到正文
Elaine Blog
返回

Kafka 客户端消息流转

分布式系统、协调与消息

Kafka 把 Topic 拆成多个只追加的分区日志。生产者决定记录进入哪个分区,Broker 按 offset 保存记录,消费组把分区分配给组内消费者。掌握这条路径,就能解释吞吐、顺序、重复和扩缩容。

Table of contents

Open Table of contents

一条记录怎样流动

Producer 序列化 key/value,按 key 或分区器选择分区,把小记录批量发送给分区 Leader。Broker 追加记录并按 acks 条件响应。Consumer 定期拉取分配给自己的分区,从当前位置读取,并提交下一次应读取的 offset。

Offset 是记录在一个分区内的位置,不是全 Topic 的全局序号。一个 Consumer Group 中,同一分区同一时刻只分配给一个成员;不同 Group 则各自读取完整数据。

key=order-42 → partition 3 → offset 901 → group=billing 的某个消费者
                              └────────→ group=audit 的某个消费者

Key 决定局部顺序

需要同一订单有序,就让同一 orderId 稳定进入同一分区。运行实验:

mvn -Dtest=DistributedMechanismsTest#equalMessageKeysStayInOnePartition test

示例用 floorMod(hash, partitionCount) 建立直觉。真实 Kafka 默认分区策略还会考虑空 key、批次和版本;增加分区数后,同一 key 的映射可能变化,因此顺序敏感 Topic 不应随意扩分区。

提交位置决定交付语义

处理前提交 offset,崩溃可能丢业务效果,接近 at-most-once;处理后提交,崩溃会重读,形成 at-least-once。常规业务采用后者,并用事件 ID 或业务唯一键幂等。

生产配置要成组检查:启用幂等生产者、合适的 acks、重试和投递超时;消费端设置合理的批量大小、处理时限与手动提交策略。监控生产错误与延迟、消费 lag、最老事件年龄、重平衡次数和重复率。

下一步

继续阅读08-16 Kafka 集群、分区、副本与再均衡,理解 Broker 故障和消费者变化时发生什么。


分享这篇文章:

上一篇
RabbitMQ 集群、Quorum Queue 与运维
下一篇
Kafka 集群、分区、副本与再均衡