Kafka 生产者与消费者
概述
生产者负责把消息高效可靠地写入 Kafka,消费者负责按语义消费。生产者关键:分区选择、批次缓冲、压缩、幂等/事务;消费者关键:消费组、Rebalance、位移提交、拉取模型。本文两边都讲透。
一、生产者架构
1.1 生产流程
Producer → 拦截器 → 序列化器 → 分区器
→ RecordAccumulator(批次缓冲)
→ Sender 线程 → 网络 → Broker
核心组件:
KafkaProducer(入口)
RecordAccumulator(缓冲批次)
Sender(发送线程)
Metadata(集群元数据)1.2 发送模型
发送方式:
同步(send().get())— 简单但慢
异步(send + 回调)— 生产常用
回调时机:
发送成功/失败后回调
可用于记录/重试补偿二、分区策略
2.1 分区选择
分区规则(未指定 key 时):
默认轮询(round-robin)→ 均匀分布
指定 key → key 哈希 → 同 key 同分区
指定分区 → 直接写指定分区
无 key 无分区器 → 粘性分区
(一批一个分区,减少小文件)2.2 自定义分区器
自定义实现:
实现 Partitioner 接口
实现 partition() 方法
注册进生产者配置
场景:
按业务 ID 路由
热点 key 打散
数据本地化2.3 顺序性保证
同分区内严格有序:
同 key 消息 → 同分区 → 顺序消费
跨分区无顺序
全局顺序方案:
单分区 + 单消费者(牺牲并行)三、批次缓冲与压缩
3.1 批次缓冲
RecordAccumulator:
按分区聚合消息到批次(RecordBatch)
达到 batch.size 或 linger.ms 触发发送
参数:
batch.size(默认 16KB)
linger.ms(默认 0)
buffer.memory(缓冲池上限)吞吐与延迟权衡:
批次越大/等待越久 → 吞吐高、延迟高
批次小/linger 小 → 延迟低、吞吐低
生产按业务取舍3.2 消息压缩
| 压缩类型 | 特点 |
|---|---|
| gzip | 压缩率高、CPU 高 |
| snappy | 均衡 |
| lz4 | 解压快 |
| zstd | 压缩率最高 |
压缩位置:
生产端压缩 → Broker 存储压缩态
→ 消费端解压
(减少网络与磁盘占用)
参数:compression.type3.3 缓冲耗尽处理
buffer.memory 满:
发送阻塞(max.block.ms)
超时抛异常
应对:
增大 buffer.memory
减小批次
提高消费能力四、幂等与事务
4.1 幂等生产者
问题:重试可能导致消息重复
方案:幂等(enable.idempotence=true)
机制:
Producer ID(PID)+ 序列号(seq)
Broker 按 PID+分区 去重
重复序列号 → 丢弃
限制:
单分区内去重
不跨会话(PID 变了不生效)4.2 事务
事务能力(跨分区原子写):
Producer 事务 + 事务协调器
流程:
开启事务 → 发送 → 提交/中止
协调器分配事务 ID
写 __transaction_state
应用:
多分区原子写
Exactly-Once 消费(读进程+写进程)4.3 幂等 vs 事务
| 能力 | 幂等 | 事务 |
|---|---|---|
| 去重范围 | 单分区 | 多分区 |
| 原子性 | 无 | 有 |
| 依赖 | 无需 | 需事务协调器 |
| 用途 | 常规生产 | 精确一次 |
建议:
默认开启幂等(开销小)
需要跨分区原子时用事务五、消费者与消费组
5.1 消费组模型
Consumer Group:
组内消费者分担分区
一个分区只能被组内一个消费者消费
不同组互不影响
消费者数 > 分区数:
多余消费者空闲(不消费)
消费者数 < 分区数:
一个消费者消费多个分区5.2 拉取模型
Pull 模型:
消费者主动拉取
控制消费速率(背压)
可暂停/继续
fetch.min.bytes/fetch.max.wait.ms与 Push 对比:
Push:服务端推送 → 消费慢会积压
Pull:客户端拉取 → 消费自控
Kafka 用 Pull,配合长轮询5.3 消费位移
位移(offset):
消费者已消费位置
存于内部 Topic:__consumer_offsets
按 group + topic + partition 记录
提交方式:
自动提交(enable.auto.commit=true,定期)
手动提交(业务确认后提交)自动提交风险:
消费了但没处理完 → 自动提交 → 重启丢数据
生产核心场景 → 手动提交
(处理成功后再提交)六、Rebalance 详解
6.1 触发时机
Rebalance(再平衡):
消费者加入/退出
消费者故障(心跳超时)
分区数变化
订阅 Topic 变化
期间:组内暂停消费,重新分配6.2 分配策略
| 策略 | 说明 |
|---|---|
| Range | 按 Topic 顺序连续分配 |
| RoundRobin | 轮询分配 |
| Sticky | 粘性,尽量保持原分配 |
Sticky 优势:
减少分区转移
降低 Rebalance 抖动
生产常用6.3 协作式 Rebalance
旧协议(Eager):
先全部撤销再分配
→ 组内短暂全部停摆
新协议(Cooperative):
增量式调整
只重分配受影响分区
→ 减少停摆Rebalance 影响:
期间消费暂停
反复 Rebalance = 消费抖动
防护:
合理的 session.timeout / heartbeat
消费处理不阻塞心跳
稳定消费组6.4 重复消费与位移丢失
风险场景:
Rebalance 前提交的位移丢失 → 重复消费
消费慢导致踢出组 → 重新分配 → 再消费
解法:
手动提交 + 提交前先暂停拉取
合理 max.poll.interval.ms
消费幂等设计七、生产消费参数表
| 场景 | 参数 | 建议 |
|---|---|---|
| 可靠写入 | acks | all |
| 防重 | enable.idempotence | true |
| 批量 | batch.size / linger.ms | 按吞吐调 |
| 压缩 | compression.type | lz4/zstd |
| 位移 | enable.auto.commit | false(核心) |
| 拉取 | max.poll.records | 按处理能力 |
| 心跳 | session.timeout.ms | 默认 45s 左右 |
八、常见问题
8.1 消费积压
原因:
消费能力 < 生产速率
分区数不足(并行度受限)
单条处理慢
处理:
增加消费者/分区
优化消费逻辑
优先追积压再恢复正常8.2 消息重复消费
原因:
Rebalance 位移丢失
手动提交时机不对
消费后提交失败
方案:
幂等消费(去重表/状态)
调整提交策略
事务(消费-处理-提交)8.3 生产超时
排查:
网络/集群压力
批次缓冲等待
acks=all 时 ISR 不足
处理:
查看请求队列/网络指标
检查 ISR 状态
调大 retries 与超时九、小结
生产端:分区策略定路由,批次压缩提性能,幂等事务保可靠;消费端:消费组定并行,拉取模型控速率,位移提交定语义。生产环境记住:可靠写入用 acks=all + 幂等,核心消费手动提交 + 幂等设计,遇到积压先看并行度再看处理逻辑。