生产、消费与可靠性
Producer 发送链路
Producer 对 Record 序列化、选择 Partition,写入批次缓冲区,再由 Sender 发送给 Partition Leader。batch.size 是批次上限目标之一,不表示只有达到该大小才发送;linger.ms 可允许短暂等待以形成更大批次。
分区策略
- 显式指定 Partition 时直接使用;
- 有 Key 时使用分区器,使相同 Key 通常进入相同 Partition;
- 无 Key 时现代客户端通常使用 Sticky Partitioning 提高批量效率,不应继续概括为简单 Round-Robin;
- 增加 Partition 后 Key 到 Partition 的映射可能变化,业务不能把默认 Hash 当作永久分片协议。
ACK 与副本约束
acks | 确认条件 | 风险 |
|---|---|---|
0 | 不等待 Broker 响应 | 客户端无法确认写入结果 |
1 | Leader 本地接收 | Leader 故障且 Follower 未复制时可能丢失 |
all | 当前 ISR 达到服务端最小同步副本约束后确认 | 可靠性最高,但不等于等待“全部副本” |
旧笔记把 acks=all/-1 描述成等待所有 Follower,这是错误的。它与 min.insync.replicas 共同决定可接受写入的最小 ISR 数量。
幂等 Producer
幂等 Producer 使用 Producer ID 和 Partition Sequence Number 去除同一 Producer 会话中的重试重复。它不能自动去重业务重复请求,也不能跨任意新 Producer 实例永久去重。
需要同时写多个 Partition 的原子性时使用事务 Producer,并设置稳定 transactional.id。事务还可把消费 Offset 与输出 Record 原子写入 Kafka,用于 consume-transform-produce。
Consumer Group
同一 Group 内,一个 Partition 在同一时刻只分配给一个 Consumer;一个 Consumer 可以消费多个 Partition。Consumer 数超过 Partition 数时,多出的 Consumer 空闲,而不是“备用节点”。
Group Coordinator 负责成员管理、分区分配和 Offset 提交。成员变化、订阅变化或故障可能触发 Rebalance。Cooperative Rebalancing 能减少一次性撤销全部分区的影响,但 Consumer 仍必须正确处理 Revoke/Assign 回调。
Offset
提交的 Offset 通常是“下一条准备消费的 Record Offset”,即已处理最后一条的 Offset + 1。Offset 保存在内部 Topic __consumer_offsets。
| 处理顺序 | 可能语义 |
|---|---|
| 先提交 Offset,再处理业务 | 故障时可能漏处理 |
| 先处理业务,再提交 Offset | 故障时可能重复处理 |
| Kafka 事务内写结果并提交 Offset | Kafka 内部端到端 Exactly-Once 路径 |
自动提交只表示客户端周期提交 Poll 返回记录的位置,不代表业务副作用一定已经完成。
消费隔离级别
事务 Producer 写入的数据可能处于未提交状态。Consumer 使用 read_committed 时只返回已提交事务数据,并跳过 Abort 事务;默认隔离级别需要按客户端配置确认。
业务幂等
Kafka 的 At-Least-Once 消费应把重复视为正常情况。常见方法包括业务唯一键、数据库唯一约束、Inbox 表、状态机条件更新和事务性 Outbox/CDC。