主题
第 4 章 Producer 深入:把一条 send() 拆成三十步
目标读者:第 3 章已经会用
kafka-console-producer.sh和confluent-kafka的produce()发消息,但被问到「acks=all到底等几个副本」「linger.ms=20和batch.size=16KB是与关系还是或关系」「幂等生产者怎么去重」就答不上来的同学。学完你会:闭着眼睛说出 Producer 内部的「序列化 → 分区 → 累加器 → Sender 线程 → in-flight → 重试 → ack 回写」全流程;能根据业务对吞吐 / 可靠性 / 延迟的诉求,主动调好 acks / retries / linger.ms / batch.size / compression / idempotence 这套组合拳;能解释幂等 Producer 的 PID + Epoch + 序列号是怎么对抗「网络抖动 + 重试」造成的重复消息。
0. 导读:你以为的 send(),和真实的 send()
新手脑海里的 Producer:
producer.send("topic", "msg") → Kafka → ✅真实的 Producer:
producer.send("topic", "msg")
→ 1) 序列化 key/value
→ 2) Partitioner 算分区
→ 3) 把消息丢到 RecordAccumulator(每分区一个双端队列)
→ 4) 立即返回(你的线程从这里继续往下走,根本没等 broker)
┊
(后台 Sender 线程独立运行) ┊
→ 5) Sender 把同一个 broker 的多个分区 batch 一次性发出去
→ 6) Broker 收到 → 按 acks 规则等副本同步 → 回 ProduceResponse
→ 7) Sender 触发回调(你的 on_delivery)
→ 8) 失败的 batch 进入重试队列,按 retry.backoff.ms 排队
→ 9) 重试时如果开了幂等,序列号让 broker 自动去重「send() 是异步的」「批量发送」「重试可能造成乱序」「幂等就是去重」——这些常识都来自上面这张图。本章把每一步拆开讲。
1. Producer 的内部架构
Producer JVM / librdkafka
┌──────────────────────────────────────────────────────────────────────┐
│ │
│ 你的业务线程 1 ──┐ │
│ 你的业务线程 2 ──┼─> produce() │
│ 你的业务线程 N ──┘ │ │
│ │ ① 序列化 + Partitioner │
│ ▼ │
│ ┌──────────────────────────────────────────────────────┐ │
│ │ RecordAccumulator │ │
│ │ 每个 (Topic, Partition) 一个 deque<Batch> │ │
│ │ │ │
│ │ topic_a-P0 ▸ [Batch(已满)] [Batch(攒中,12KB)] │ │
│ │ topic_a-P1 ▸ [Batch(攒中, 3KB)] │ │
│ │ topic_b-P0 ▸ [Batch(攒中, 8KB)] │ │
│ └──────────────────────────────────────────────────────┘ │
│ │ ② Sender 后台线程定期检查 │
│ ▼ │
│ ┌──────────────────────────────────────────────────────┐ │
│ │ Sender Thread │ │
│ │ - 把发往同一 broker 的多 partition batch 合并成 │ │
│ │ 一个 ProduceRequest(按 broker 维度) │ │
│ │ - 维护 in-flight 请求队列(受 max.in.flight 限制) │ │
│ │ - 失败的 batch 放回 accumulator 等下一轮 │ │
│ └──────────────────────────────────────────────────────┘ │
│ │ ③ 网络 I/O │
│ ▼ │
│ ┌──────────────┐ │
│ │ Broker │ ack 后回调 on_delivery │
│ └──────────────┘ │
└──────────────────────────────────────────────────────────────────────┘🍱 生活类比:Producer 像一个集装箱码头——
- 你的业务线程 = 货车司机,把包裹(消息)卸到对应的集装箱(Partition)。
- RecordAccumulator = 集装箱仓库,按目的地分门别类堆放。
- Sender 线程 = 起重机,定时把同一艘船(Broker)要装的多个集装箱一次性吊上船。
linger.ms= 「再等多 5 分钟,看还有没有同方向的货」;batch.size= 「单个集装箱最多塞 16KB」;任何一个先满足就发车。
2. 同步 vs 异步 send(confluent-kafka 篇)
2.1 异步 send(默认/推荐)
python
from confluent_kafka import Producer
p = Producer({"bootstrap.servers": "127.0.0.1:9092", "linger.ms": 5})
def cb(err, msg):
if err: print("FAIL:", err)
else: print(f"OK: {msg.partition()}@{msg.offset()}")
for i in range(1000):
p.produce("learn.04.demo", value=f"v{i}".encode(), on_delivery=cb)
p.poll(0) # 非阻塞,触发已就绪回调
p.flush(10) # 阻塞,等所有未完成消息完结 + 触发回调特征:
produce()不阻塞(除非本地 buffer 满,会抛BufferError)。- 真正的 ack 通过
on_delivery回调通知。 - 单线程能跑 50w-100w msg/s。
2.2 同步 send(一条等一条)
confluent-kafka-python 没有原生同步 API,自己 wrap:
python
import threading
def send_sync(producer, topic, value, key=None, timeout=10):
done = threading.Event()
result = {}
def cb(err, msg):
result["err"] = err
result["msg"] = msg
done.set()
producer.produce(topic, value=value, key=key, on_delivery=cb)
producer.poll(0)
if not done.wait(timeout):
raise TimeoutError("send timeout")
if result["err"]:
raise Exception(result["err"])
return result["msg"]
msg = send_sync(p, "learn.04.demo", b"hello")
print(msg.partition(), msg.offset())特征:
- 一条阻塞等一条,吞吐瞬间降到几百到几千 msg/s(受 RTT + acks 影响)。
- 适用于「关键消息必须 ack 后才继续」场景,例如订单流水写完才允许扣款。
2.3 性能对比
测试机:单 Broker KRaft,本地 loopback,单分区,1KB 消息,acks=all,无压缩。
| 模式 | 吞吐 (msg/s) | p99 延迟 (ms) | CPU |
|---|---|---|---|
| 异步 + linger=5ms | 168,000 | 18 | 35% |
| 异步 + linger=0ms | 92,000 | 9 | 50% |
| 同步(每条 wait) | 1,200 | 0.8 | 8% |
结论:没有特殊业务原因就用异步 + 回调,再用 flush() 把控边界。
3. acks=0/1/-1:Kafka 的「可靠性三档」
acks 决定 Producer 等几个副本写完才认为消息「成功」。
3.1 acks=0:发了就忘
- 吞吐最高(不等任何 ack)。
- 可能丢消息:Broker 没收到 / 收到但写失败 / Leader 切换都丢。
- 零重试(没有响应就没有重试触发条件)。
- 使用场景:可丢失的指标 / 日志(例如 metrics 上报、A/B 实验事件采样)。
3.2 acks=1:Leader 写完就 ack
- 吞吐次高(只等 Leader 落盘)。
- 会丢消息:Leader 写完 ack,但 Follower 还没拉到 → Leader 宕机 → 消息丢失。
- 默认值(旧版本,3.0 起
acks默认all,但很多客户端依然显式配置acks=1)。 - 使用场景:能容忍极少量丢失(每年丢千分之一可接受),且追求高吞吐的业务,例如内部日志总线。
3.3 acks=-1(all):所有 ISR 都写完才 ack
- 吞吐最低(要等所有 ISR 副本同步)。
- 不丢消息(前提:min.insync.replicas ≥ 2):哪怕 Leader 立刻宕机,剩下任意 ISR 副本接任都能保住消息。
- 必须搭配
min.insync.replicas:min.insync.replicas=1+acks=all仍然可能丢(ISR 只剩 Leader 时退化为 acks=1)。- 业界推荐:RF=3 + min.insync.replicas=2 + acks=all。
- 使用场景:金融、订单、计费、Exactly Once 必须的前置条件。
3.4 三档对照表
| acks | 吞吐 | 可能丢消息? | 典型 p99 (ms) | 重试有意义? | 使用场景 |
|---|---|---|---|---|---|
| 0 | 极高 | ✅ 会 | 1-3 | ❌ 没响应也不重试 | 指标采样 |
| 1 | 高 | ⚠️ 极小概率 | 5-10 | ✅ | 一般业务日志 |
| -1 | 中 | ❌ 不丢(配 min.insync ≥ 2) | 15-30 | ✅ | 金融、订单、EOS |
🍱 生活类比:寄信
acks=0= 把信扔进邮筒就走,邮局收没收到不知道。acks=1= 邮局柜员盖了戳就给你回执,至于他们内部归档没归档不管。acks=all= 邮局必须把信送到收件方手上、对方签收回执,你才认为寄成功。
4. retries / delivery.timeout.ms / max.in.flight 的三角关系
这三个参数搅在一起,是「重试 + 顺序 + 延迟」的核心。
4.1 三个参数各自含义
| 参数 | 默认值 | 含义 |
|---|---|---|
retries | Integer.MAX_VALUE(3.0+) | 单个 batch 的最大重试次数;只决定「最多试几次」 |
retry.backoff.ms | 100 | 两次重试之间的最小间隔 |
delivery.timeout.ms | 120000 | 从 produce() 入队到 ack(含全部重试)的总超时 |
request.timeout.ms | 30000 | 单次 ProduceRequest 等响应的超时 |
max.in.flight.requests.per.connection | 5 | 单个 broker 连接上未收到响应的请求数上限 |
4.2 它们怎么交互
produce() 入队 ─────────────────────────────────────►
delivery.timeout.ms(兜底总闸门)
尝试 1 ─request.timeout.ms─► [失败]
└ retry.backoff.ms ─►
尝试 2 ─request.timeout.ms─► [失败]
└ retry.backoff.ms ─►
尝试 3 ─request.timeout.ms─► [成功 ✅]
任何时候只要:
- 重试次数超过 retries
- 总耗时超过 delivery.timeout.ms
都会抛错给 on_delivery。关键约束:
delivery.timeout.ms ≥ linger.ms + request.timeout.ms × (retries + 1)
+ retry.backoff.ms × retries如果 delivery.timeout.ms 设得太小,重试次数白配。
4.3 max.in.flight 与「重试乱序」
不开幂等时,max.in.flight > 1 会出现:
Producer 发出:
Request A (batch1=msg1,msg2) → in-flight
Request B (batch2=msg3,msg4) → in-flight
Broker 端结果:
A 失败需重试
B 成功,先落盘
Producer 重发 A:
msg1, msg2 落盘在 msg3, msg4 之后 ❌ 顺序错乱经典规则:
| 场景 | 推荐配置 |
|---|---|
| 普通业务,能接受少量乱序 | max.in.flight=5(默认),不开幂等 |
| 需要分区内严格顺序,不开幂等 | max.in.flight=1,吞吐降到 1/5 |
| 需要严格顺序 + 不丢 + 不重 | enable.idempotence=true,broker 内部按序列号重排,可以放心 max.in.flight=5(≤5) |
⚠️ 开了幂等 (enable.idempotence=true) 时,client 会自动把 max.in.flight 限制到 ≤ 5、acks=all、retries>0,不需要你手动配。
5. linger.ms / batch.size / compression:吞吐三剑客
5.1 它们的关系
linger.ms 和 batch.size 是 「或」关系:任意一个先满足就发车。
batch 当前大小 ≥ batch.size ──► 立即发出
linger.ms 超时 ──► 立即发出
哪个先到,触发哪个5.2 为什么批次能省时间?
每条消息单独发:
[1 RTT 网络] + [1 次磁盘写] + [1 次副本同步] = ~10ms / 条
1000 条 → 10000ms = 10s批次发:
[1 RTT 网络] + [1 次磁盘顺序写] + [1 次副本同步] = ~12ms / 整批 1000 条
吞吐 ~80x⚠️ 这就是「Kafka 高吞吐」的核心机制之一——用 batch 把每条消息的固定开销摊薄到 1/N。
5.3 实测数据:linger 的威力
单 Broker KRaft,6 分区,1KB 消息,acks=all,单 Producer,无压缩。
| linger.ms | batch.size | 吞吐 (MB/s) | 平均延迟 (ms) | p99 延迟 (ms) | 网络包数 |
|---|---|---|---|---|---|
| 0 | 16384 | 22 | 5 | 12 | 极多(每条一包) |
| 5 | 16384 | 89 | 8 | 18 | 中 |
| 20 | 16384 | 154 | 23 | 40 | 少 |
| 100 | 16384 | 168 | 105 | 130 | 极少 |
| 5 | 65536 | 132 | 9 | 22 | 中 |
| 5 | 1MB | 154 | 11 | 25 | 少 |
观察结论:
linger.ms=0即「来一条发一条」,吞吐惨不忍睹但延迟最小。linger.ms在 5-20 之间是大多数业务的甜蜜点:吞吐 6-7 倍,延迟只多个十几毫秒。batch.size单方面调大,效果不如linger.ms明显(因为没消息就攒不起来)。linger.ms ≥ 100后吞吐基本到顶,再加只是多攒延迟。
5.4 compression:CPU 换网络/磁盘
支持 5 种:none / gzip / snappy / lz4 / zstd。
同上环境,1KB JSON(重复字段多,可压缩性强),acks=all,linger.ms=20。
| compression | 压缩率 | Producer 吞吐 (MB/s) | Broker 入站带宽 (MB/s) | Producer CPU |
|---|---|---|---|---|
| none | 1.0× | 154 | 154 | 28% |
| gzip | 6.8× | 92 | 14 | 88% |
| snappy | 3.1× | 138 | 45 | 41% |
| lz4 | 3.0× | 162 | 54 | 35% |
| zstd | 5.5× | 156 | 28 | 52% |
经验法则:
- 2025 年生产环境首选
zstd:压缩率接近 gzip,CPU 是 gzip 的一半。 - 极致低 CPU →
lz4。 - 极致省存储 →
gzip(但 CPU 痛苦)。 - 不可压数据(图片、加密内容)→
none,硬压只是浪费 CPU。
⚠️ 压缩在 batch 级:单条消息开 compression=zstd 没意义,必须批起来才有压缩比。所以压缩 + linger.ms 是互补:linger 攒大批 → 压缩比更高 → 网络更省。
5.5 「优化前 / 优化后」对照(综合)
业务:每条消息 800 字节 JSON 订单,要求不丢消息,吞吐越高越好。
| 配置 | 吞吐 (msg/s) | p99 延迟 (ms) | 备注 |
|---|---|---|---|
| 优化前:默认 acks=1, linger.ms=0, compression=none | 91,000 | 8 | 默认配置 |
| 改 acks=all (RF=3, min.isr=2) | 38,000 | 22 | 可靠性升一级,吞吐降一半 |
| 加 linger.ms=20 | 142,000 | 35 | 攒批起作用 |
| 加 compression=zstd | 156,000 | 38 | 网络省 5x |
| 开 enable.idempotence | 152,000 | 40 | 吞吐略降,去重保障 |
| 优化后:acks=all + linger=20 + zstd + idempotence | 152,000 | 40 | 吞吐 1.7x,且不丢不重 |
6. Partitioner:消息怎么落到分区
6.1 三种内置策略
| Partitioner | 行为 | 适用 |
|---|---|---|
| DefaultPartitioner (3.0+) | Key 存在 → Murmur2(key) % numPartitions;Key 为 null → Sticky | 通用 |
| StickyPartitioner (老版本,2.4-2.x) | Key=null 时连续多条粘在同一个分区,攒批更快 | 老版本默认 |
| UniformStickyPartitioner (2.4+) | Key=null 时按 batch 切换分区,更均匀 | 流量较大 |
| RoundRobinPartitioner | 严格轮询 | 极度需要均匀分布的小流量 |
6.2 Key 存在时的 Hash
partition = Murmur2(key_bytes) & 0x7FFFFFFF % numPartitions- 同一个 Key 永远落到同一个分区(除非分区数变了)。
- 这就是「单分区有序,跨分区无序」的来源。
__consumer_offsets内部使用(group, topic, partition)作为 Key,保证同一消费组的 offset 总在同一个分区,方便消费协调。
6.3 Key=null 时为什么要 Sticky?
老版本 Key=null 用 round robin → 每条消息进不同分区 → 每个分区都攒不起 batch → 网络包多、压缩比低。
StickyPartitioner:
- 一段时间内 粘在同一个分区,让 batch 攒满或超 linger 后再切换。
- 实测吞吐提升 30%-50%,延迟下降 20%。
⚠️ 副作用:短时间内消息分布不均,长时间统计才均匀。
6.4 自定义 Partitioner
confluent-kafka-python 通过 partitioner 配置或 produce(... partition=N) 显式指定分区。Java 版要实现 Partitioner 接口。
业务案例:「VIP 订单走 0 号分区,普通订单按 user_id 哈希」:
python
def my_partitioner(key, value, all_partitions, available_partitions):
if value and b'"vip":true' in value:
return 0
if key:
import mmh3
return (mmh3.hash(key) & 0x7FFFFFFF) % len(all_partitions)
return None # 让默认策略处理
# librdkafka 不支持 Python 回调式 partitioner,要显式指定 partition:
def send(producer, topic, key, value):
p = my_partitioner(key, value, [0,1,2,3,4,5], [0,1,2,3,4,5])
if p is not None:
producer.produce(topic, key=key, value=value, partition=p)
else:
producer.produce(topic, key=key, value=value)7. 幂等生产者:Exactly-Once 的入门
7.1 重复消息从哪儿来?
不开幂等时,重试就是重复的重要来源。
7.2 幂等的两个核心:PID + 序列号
开 enable.idempotence=true 后:
- Producer 启动时,向某个 Broker 发
InitProducerIdRequest,broker 给它分配一个全局唯一的 ProducerId(PID)+ Epoch。 - 每个 (PID, partition) 维护一个自增序列号 sequence number,每条消息从 0 递增。
- Broker 端为每个 (PID, partition) 缓存「最后 5 个序列号」(受
max.in.flight ≤ 5限制)。
Broker 收到消息时的判断:
incoming_seq == last_seq + 1 → 接受,落盘,更新 last_seq
incoming_seq <= last_seq → 丢弃(重复消息),返回成功
incoming_seq > last_seq + 1 → 拒绝(OutOfOrderSequence 异常),可能是丢了消息7.3 ASCII 详图
Producer 启动:
┌──────────────────────────────────────────────────┐
│ InitProducerIdRequest │
│ ↓ │
│ Broker (Transaction Coordinator) │
│ - 分配 PID = 1001, Epoch = 0 │
│ - 持久化到 __transaction_state 内部 topic │
│ ↓ │
│ Producer 拿到 PID=1001, Epoch=0 │
└──────────────────────────────────────────────────┘
发消息(partition=2):
┌──────────────┐ ┌──────────────┐
│ Producer │ │ Broker │
│ PID=1001 │ │ partition 2 │
│ Epoch=0 │ │ last_seq=-1 │
│ seq[2]=0 │ │ │
└──────┬───────┘ └──────┬───────┘
│ │
│ msg(seq=0) ──────────────────► │ seq=0 == last+1 ✅
│ ◄────────── ack │ last_seq=0
│ │
│ msg(seq=1) ──────────────────► │ seq=1 == last+1 ✅
│ (响应丢失) │ last_seq=1
│ │
│ msg(seq=1) [重试] ──────────► │ seq=1 <= last_seq=1 → 丢弃,回 OK
│ ◄────────── ack │ last_seq=1(不变)
│ │
│ msg(seq=2) ──────────────────► │ seq=2 == last+1 ✅
│ ◄────────── ack │ last_seq=2效果:
- ✅ 重试不会造成重复(broker 自动去重)。
- ✅ 顺序保证(即使 max.in.flight=5,broker 收到乱序也会按 seq 重排)。
- ⚠️ 只在「单 Producer 会话内」生效。Producer 重启 → 拿到新 PID → broker 视为新生产者 → 老消息可能被去重,新消息无法被识别为「老 PID 的重发」。
- ⚠️ 跨分区不保证:PID + seq 是按 partition 维护的,跨分区的原子性要靠事务 Producer(第 13 章)。
7.4 Epoch 解决「僵尸 Producer」
如果一个老 Producer 网络卡住,新 Producer 已经接管(场景:失败重启),broker 通过 Epoch 拒绝老 PID 的消息:
Broker 维护 (PID, max_epoch)
老 Producer (PID=1001, Epoch=0) 的消息 → broker 看 max_epoch=1 → 拒绝
新 Producer (PID=1001, Epoch=1) 的消息 → 接受事务 Producer 的 transactional.id 复用就是基于这个机制(同 transactional.id → 同 PID → 自增 Epoch)。
7.5 开启幂等的代价
| 项 | 影响 |
|---|---|
| 吞吐 | 几乎无影响(< 5%) |
| 内存 | broker 端每分区多缓存 5 个序列号 |
| 强制配置 | acks=all、max.in.flight ≤ 5、retries > 0(违反会启动失败) |
| Producer 重启 | 新 Session 拿不到老 PID,去重边界重置 |
结论:3.0+ 开 enable.idempotence=true 几乎是免费午餐,建议默认开。
7.6 Python 示例
python
from confluent_kafka import Producer
p = Producer({
"bootstrap.servers": "127.0.0.1:9092",
"enable.idempotence": True,
# 下面三个开了幂等会自动设置,不写也行;写出来是为了教学
"acks": "all",
"retries": 1000000,
"max.in.flight.requests.per.connection": 5,
})
for i in range(100):
p.produce("learn.04.idempotent", key=f"k{i%5}".encode(), value=f"v{i}".encode())
p.flush(10)配套脚本
code/idempotent_producer.py演示「故意 kill 网络重试 → broker 去重」的完整链路。
8. 错误处理与重试策略
8.1 错误分类
confluent-kafka 把错误按「是否可重试」分两类:
| 类别 | 例子 | 处理 |
|---|---|---|
| 可重试 | LEADER_NOT_AVAILABLE、NOT_LEADER_FOR_PARTITION、NETWORK_EXCEPTION、REQUEST_TIMED_OUT | 自动重试(受 retries / delivery.timeout 限制) |
| 不可重试 | RECORD_TOO_LARGE、INVALID_TOPIC、UNKNOWN_TOPIC_OR_PARTITION(持续)、AUTHORIZATION_FAILED | 立即 fail,回调里收到 |
8.2 BufferError 的处理
produce() 在本地队列满时抛 BufferError:
python
import time
from confluent_kafka import Producer, KafkaException
def safe_produce(p, topic, value, retries=5):
for _ in range(retries):
try:
p.produce(topic, value=value)
p.poll(0)
return
except BufferError:
p.poll(0.5) # 让回调释放队列
raise RuntimeError("队列长期满,放弃")8.3 死信处理
业务无法处理(消息格式错、不可重试错误)时,常见做法:
python
def cb(err, msg):
if err:
# 写到死信 Topic
dlq_producer.produce(
"learn.04.dlq",
key=msg.key(),
value=msg.value(),
headers={"original-topic": msg.topic(), "error": str(err)},
)第 16 章 Kafka Connect 的 DLQ 机制是更体系化的方案。
9. 与其他 MQ 对比
| 维度 | Kafka Producer | RabbitMQ Producer | RocketMQ Producer |
|---|---|---|---|
| 写入模型 | 批量异步(每分区 batch + Sender 线程) | 单条 publish + confirm | 同步 / 异步 / OneWay 三模式 |
| 顺序保证 | 单分区严格有序(开幂等才保险) | 单队列有序 | MessageQueueSelector + 同步发送可保证 |
| Exactly Once | 幂等 (PID+seq) + 事务 (transactional.id) | RabbitMQ 5+ 才有 streams 支持 | 半消息事务 |
| 压缩 | batch 级 (gzip/snappy/lz4/zstd) | 单条 / 不内置压缩 | 单条 zip |
| 路由 | 按 Key Hash | Exchange + Routing Key 灵活 | Tag + 自定义 Selector |
| 吞吐量级 | 单 Producer 50w-100w msg/s | 单 connection 1w-10w msg/s | 单 Producer 数十万 |
📌 关键差异:Kafka Producer 是「为吞吐优化」的,本质是「批 + 异步 + 顺序写」。如果业务要求每条消息都立即 RPC 风格 ack(如银行转账),Kafka Producer 的同步模式吞吐还不如直接用 gRPC,应该考虑别的方案或上事务。
10. 调优清单
Producer 调优 6 步走
1) 业务先定 SLO:吞吐 / 延迟 / 可靠性,三选二
2) 选 acks(0/1/all)
- 不丢消息 → all + min.insync.replicas≥2 + RF≥3
3) 开 enable.idempotence=true(几乎是免费的)
4) 攒批
- linger.ms = 5~50(看延迟容忍度)
- batch.size = 16384~1048576
5) 压缩
- 默认开 zstd / lz4
6) 监控关键指标
- record-send-rate (吞吐)
- record-error-rate (错误率)
- record-retry-rate (重试率,高就要排查 broker)
- record-queue-time-avg/p99 (在 accumulator 里待了多久)
- request-latency-avg/p99 (broker 处理 + 网络)
- batch-size-avg (平均批次大小)11. 实战脚本一览
本章配套目录
04_producer/:
init.sh—— 创建本章用到的 Topic(learn.04.demo、learn.04.idempotent、learn.04.bench等)。code/sync_producer.py—— 自己 wrap 同步发送,演示「等 ack 再继续」。code/async_producer.py—— 标准异步 + 回调 + flush 边界。code/idempotent_producer.py—— 开启幂等,演示 PID 与序列号的运作。code/benchmark_compression.py—— 可复现压测:4 种压缩算法 × 3 种 linger 配置的吞吐 / 网络对比。demo.html—— 浏览器交互:滑块调 linger / batch / compression / acks,实时仪表盘 + 批次聚合动画。
12. 小结
Producer 深入
├─ 内部架构
│ ├─ 业务线程 → 序列化 → Partitioner → RecordAccumulator
│ └─ Sender 后台线程 → 合并 → 网络 → 重试 → 回调
├─ 同步 vs 异步
│ ├─ confluent-kafka 默认异步
│ └─ 同步要自己 wrap,吞吐降 100x
├─ acks 三档
│ ├─ 0 = 不等,可能丢
│ ├─ 1 = Leader 确认,仍可能丢
│ └─ all = ISR 全确认(配 min.insync ≥ 2)
├─ 重试三角
│ ├─ retries / delivery.timeout.ms / max.in.flight
│ └─ 开幂等可放心 max.in.flight ≤ 5 不乱序
├─ 吞吐三剑客
│ ├─ linger.ms (甜蜜点 5-20ms)
│ ├─ batch.size(16KB-1MB)
│ └─ compression(首选 zstd)
├─ Partitioner
│ ├─ 有 Key → Murmur2 Hash
│ └─ 无 Key → Sticky 攒批友好
└─ 幂等生产者
├─ PID + Epoch + 序列号
├─ Broker 按 (PID, partition) 缓存最近 5 个 seq
└─ 单 Session 内不丢不重,跨 Session 要事务13. 面试高频题(6 题)
Q1:acks=0/1/-1 各自的语义、典型丢消息场景、和性能影响?
考察点:可靠性核心,最高频题之一。
答案:
- acks=0:Producer 发出请求后不等任何响应就认为成功。
- 可能丢:Broker 没收到 / 收到但写失败 / Leader 切换全丢。
- 吞吐最高,但重试无意义(没响应就没失败信号)。
- 适用:可丢的指标采样、A/B 实验事件。
- acks=1:Leader 写入本地 log 就回 ack。
- 仍可能丢:Leader 写完但 Follower 还没拉到 → Leader 宕机 → 新 Leader 没这条消息。
- 适用:内部日志、能容忍极少量丢失。
- acks=-1 / acks=all:等 ISR 中所有副本都同步完才 ack。
- 不丢的前提是
min.insync.replicas ≥ 2,否则 ISR 只剩 Leader 时退化为 acks=1。 - 业界推荐组合:RF=3 + min.insync.replicas=2 + acks=all。
- 适用:金融、订单、事务、Exactly Once。
- 不丢的前提是
- 吞吐对比(同环境):acks=0 ≈ 250k msg/s,acks=1 ≈ 170k,acks=all ≈ 90k。
- Kafka 3.0 起默认值改成 all(之前是 1),就是为了「默认安全」。
加分项:能解释 unclean.leader.election.enable=true 在 ISR 全挂时允许非 ISR 副本上位,会让 acks=all 也丢消息(已经 ack 的消息被截断);提到 __transaction_state 和 __consumer_offsets 内部 topic 的 min.insync.replicas 是 2,正是基于这个原理。
Q2:max.in.flight.requests.per.connection 大于 1 一定会乱序吗?开了幂等还会吗?
考察点:重试 / 顺序 / 幂等的交互,深度题。
答案:
- 不开幂等:
max.in.flight > 1时,多个请求并发在网络上,一旦前一个失败重试,可能比后一个晚到 → 同一分区内消息乱序。- 严格保序的传统做法:
max.in.flight=1,吞吐降到 1/5。
- 严格保序的传统做法:
- 开幂等 (
enable.idempotence=true):- Broker 维护 (PID, partition) 维度的
last_seq,每条消息带自增 sequence number。 - 即使 client 把 5 个请求并发发出,broker 也会按 seq 顺序检查:seq != last+1 时拒绝。
- 失败重试时,broker 看到老 seq 直接丢弃(去重),不影响顺序。
- client 端会自动把 max.in.flight 限制在 ≤ 5(库内部强制,超过会启动失败)。
- Broker 维护 (PID, partition) 维度的
- 结论:开幂等 + max.in.flight ≤ 5 是「既保序又高吞吐」的最优组合,不开幂等的话只能在「max.in.flight=1(保序)」和「>1(高吞吐但乱序)」之间二选一。
- 跨分区:同一 Producer 写多个分区时,Kafka 不保证跨分区顺序——这是分区模型的本质,不可改变。需要全局顺序就只能用单分区(牺牲并行度)。
加分项:提到事务 Producer(transactional.id)让多个分区的写入原子化,但顺序仍是分区内的;Kafka Streams 内部就是用「按 Key 路由保证 partition 内顺序」实现一致性。
Q3:linger.ms / batch.size 是「与」还是「或」?怎么调?
考察点:吞吐调优基本功。
答案:
- 「或」关系:任意一个先满足就触发发送。
batch 当前大小 ≥ batch.size→ 发车。linger.ms计时器超时 → 发车。
linger.ms的取舍:0:来一条发一条,吞吐惨(每条消息都付出固定开销),延迟最低。5-20:大多数业务的甜蜜点,吞吐能提 5-10x,延迟只多十几毫秒。>50:吞吐基本到顶,再加只是多攒延迟。
batch.size的取舍:- 默认
16384(16KB)。 - 调大(64KB-1MB)能让单批包含更多消息,配合 linger 更高效;超过 1MB 通常没收益(被 broker 端
message.max.bytes限制)。
- 默认
- 「调多少」实操:
- 业务延迟容忍 50ms 以内 → linger=10,batch.size=64KB,吞吐已经能跑满千兆网。
- 能容忍 200ms+ 的离线管道 → linger=100,batch.size=1MB,吞吐到顶。
- 配套:开
compression=zstd才能让大 batch 的网络收益最大化(小 batch 压缩比低)。 - 指标:盯
batch-size-avg,如果远小于batch.size,说明流量太稀疏,调大 batch 没用,应该调大linger.ms或合并业务。
加分项:提到 librdkafka 还有 queue.buffering.max.ms 是 linger.ms 的别名;Java 中 KafkaProducer 的 metrics() 能拿到 record-queue-time-avg 直观看到消息在 accumulator 里待了多久。
Q4:enable.idempotence=true 的实现原理?能不能跨 Producer 重启去重?
考察点:幂等的内部机制。
答案:
- 核心三件套:
- PID(Producer ID):Producer 启动时向 broker 申请,全局唯一。
- Epoch:每次 PID 被「重新认领」(事务场景下复用 transactional.id 时)就 +1,旧 Epoch 的请求被拒。
- 序列号 (sequence number):按 (PID, partition) 维度从 0 自增,每条消息带一个。
- broker 端去重:
- 每个 (PID, partition) 缓存最近 5 个 seq(受 max.in.flight ≤ 5 限制)。
- 收到
seq == last_seq + 1才接受,seq <= last_seq视为重复直接丢弃返回成功,seq > last_seq + 1视为消息丢失返回OutOfOrderSequenceException。
- 解决的问题:网络抖动导致 client 重试 → broker 去重 → 单 Session 内不重不丢。
- 跨 Session 不行:
- Producer 重启会重新
InitProducerIdRequest,拿到新 PID。 - broker 视为完全不同的 Producer,去重历史失效。
- 想跨重启幂等:用事务 Producer,配
transactional.id,broker 会让同一 transactional.id 复用 PID 并自增 Epoch,把旧 Epoch 的「僵尸消息」屏蔽。
- Producer 重启会重新
- 强制约束:开了 idempotence 会自动强制:
acks=allmax.in.flight.requests.per.connection ≤ 5retries > 0- 违反会在客户端启动时报
ConfigException。
- 代价:吞吐降低 < 5%,broker 内存多缓存几个 seq,几乎免费,3.0+ 强烈建议默认开。
加分项:提到 Kafka 3.0 起 enable.idempotence=true 是默认值;__transaction_state 内部 topic 存储事务状态,去重的 PID/Epoch 信息存在 broker 内存(重启会从 __transaction_state 恢复事务上下文,但单纯幂等的 PID 不持久化)。
Q5:delivery.timeout.ms、request.timeout.ms、retries、retry.backoff.ms 怎么配合?
考察点:超时 + 重试模型。
答案:
- 各自含义:
request.timeout.ms:单次 ProduceRequest 等响应的超时(默认 30s)。retry.backoff.ms:两次重试之间的最小间隔(默认 100ms)。retries:单个 batch 的最大重试次数(3.0+ 默认Integer.MAX_VALUE,由delivery.timeout.ms兜底)。delivery.timeout.ms:从 produce() 入队到最终 ack 的总超时(默认 120s),是兜底总闸门。
- 关系约束:设小了
delivery.timeout.ms ≥ linger.ms + request.timeout.ms × (retries + 1) + retry.backoff.ms × retriesdelivery.timeout.ms会让 retries 失效。 - 触发失败的条件:先到先触发——
- 重试次数超过
retries。 - 总耗时超过
delivery.timeout.ms。
- 重试次数超过
- 3.0 起的最佳实践:
retries不用动(默认极大),重点调delivery.timeout.ms控制业务能接受的最长等待。 - 拓展:
linger.ms也算入delivery.timeout.ms的总时间;socket.connection.setup.timeout.ms是建立连接的超时,与 request.timeout 区分。
加分项:提到 broker 端的 replica.lag.time.max.ms 控制 ISR 收缩,与 Producer 端 delivery.timeout.ms 配合决定「Leader 切换时 Producer 等多久」。
Q6:Producer 内部 RecordAccumulator 是什么?如何避免 BufferError?
考察点:内部数据结构、内存控制。
答案:
- RecordAccumulator 是 Producer 在客户端进程内的内存缓冲区:
- 按 (Topic, Partition) 维度组织,每个分区一个双端队列
Deque<ProducerBatch>。 - 业务线程
produce()把消息追加到对应分区的最后一个 batch,满了就新建一个 batch。 - Sender 后台线程从这里 drain(按 broker 维度合并)发出。
- 按 (Topic, Partition) 维度组织,每个分区一个双端队列
- 总大小由
buffer.memory(默认 32MB)控制,所有分区的 batch 共享。 - BufferError / TimeoutException 的成因:业务发送速度 > Sender 发送速度 + broker 处理速度,buffer 被打满。
- 解决思路:
- 加大 buffer:
buffer.memory=128MB(注意会占用 client 内存)。 - 加快发送:调大
max.in.flight、开压缩、加 broker 数量、扩分区。 - 业务限速:在业务侧加令牌桶,避免突发超过下游能力。
- 改异步队列:业务线程
try/except BufferError,sleep 后重试,让 Sender 喘息。
- 加大 buffer:
- librdkafka 的 Python 客户端:
queue.buffering.max.messages(默认 100k 条)、queue.buffering.max.kbytes(默认 1GB)控制本地队列;produce()抛BufferError时调producer.poll(1)触发回调释放空间。 - 监控:
record-queue-time-avg/p99看消息在 accumulator 里待了多久;buffer-available-bytes看剩余空间;record-send-rate看实际发送速度。
加分项:提到 Kafka Streams 内部其实就是封装了一个状态机驱动的 Producer + Consumer,accumulator 的容量配置在高并发流处理里要重点调;事务 Producer 会让 accumulator 在 commit 前保留所有未提交 batch,可能放大内存占用。
下一章
05_consumer.md会从「读端」反方向把链路再走一遍:subscribe vs assign、poll 模型、Offset 提交、Coordinator 协议、Rebalance、seek/pause/resume——把 Consumer 的所有细节也拍扁。
🎬 可视化演示
演示加载缓慢或样式异常?点此在新标签页打开 ↗
💻 示例代码
python
"""
第 4 章 - 异步发送 + 回调 + flush 边界
=================================
标准的高吞吐姿势:
1) produce() 异步入队;
2) 每次 produce 后 poll(0) 触发已就绪回调;
3) 业务边界处 flush() 阻塞等所有未完成消息出去;
4) BufferError 时退避并继续。
运行:
bash ../init.sh
python async_producer.py # 默认发 10000 条
python async_producer.py 100000
"""
from __future__ import annotations
import sys
import time
from confluent_kafka import Producer
BOOTSTRAP = sys.argv[2] if len(sys.argv) > 2 else "127.0.0.1:9092"
N = int(sys.argv[1]) if len(sys.argv) > 1 else 10000
TOPIC = "learn.04.demo"
success = [0]
failure = [0]
sample = []
def on_delivery(err, msg):
if err is not None:
failure[0] += 1
if failure[0] <= 3:
print(f" ❌ {err}")
else:
success[0] += 1
if len(sample) < 5:
sample.append(
f"P{msg.partition()}@{msg.offset()} key={msg.key()} val_size={len(msg.value())}"
)
def main() -> None:
p = Producer(
{
"bootstrap.servers": BOOTSTRAP,
"client.id": "ch4-async-producer",
"acks": "all",
"enable.idempotence": True,
"compression.type": "zstd",
"linger.ms": 10,
"batch.size": 65536,
# 单连接最多 5 个未 ack 请求(开了幂等的安全上限)
"max.in.flight.requests.per.connection": 5,
"queue.buffering.max.messages": 200000,
}
)
print(f"🚀 异步发送 {N} 条到 {TOPIC}(acks=all + idempotence + zstd + linger=10)")
payload = ("x" * 800).encode()
start = time.perf_counter()
for i in range(N):
while True:
try:
p.produce(
TOPIC,
key=f"k{i % 100}".encode(),
value=payload + f"-{i}".encode(),
on_delivery=on_delivery,
)
break
except BufferError:
# 队列满,让回调释放空间
p.poll(0.5)
# 每 1000 条触发一次回调,避免回调队列堆积
if i % 1000 == 0:
p.poll(0)
remaining = p.flush(timeout=30)
elapsed = time.perf_counter() - start
print("\n📊 结果")
print(f" 成功 {success[0]} / 失败 {failure[0]} / 未发出 {remaining}")
print(f" 总耗时 {elapsed:.2f}s,吞吐 {success[0] / elapsed:.0f} msg/s")
bytes_total = success[0] * (len(payload) + 5)
print(f" 数据量 {bytes_total / 1024 / 1024:.1f} MB,{bytes_total / elapsed / 1024 / 1024:.1f} MB/s")
print("\n 样本:")
for s in sample:
print(f" - {s}")
if __name__ == "__main__":
main()python
"""
第 4 章 - 压缩算法对比压测
=================================
对比 none / gzip / snappy / lz4 / zstd 五种压缩在相同数据量下的:
- Producer 端实际吞吐 (msg/s, MB/s)
- 平均 batch 大小(反映压缩率)
- Producer CPU 时间(user + sys)
- 总耗时
数据:每条 ~1KB 的 JSON 风格 payload(字段重复多,可压缩性强)。
默认 N = 100,000 条,可通过参数调整。
运行:
bash ../init.sh
python benchmark_compression.py # 默认 100k
python benchmark_compression.py 500000 # 50w
python benchmark_compression.py 100000 myhost:9092
输出(示例,单 broker 单机 KRaft):
| codec | msgs/s | MB/s | avg_batch_kb | cpu_s |
| none | 152340 | 145.2 | 86.3 | 4.21 |
| gzip | 91230 | 87.0 | 12.7 | 9.55 |
| snappy | 138900 | 132.4 | 28.4 | 5.18 |
| lz4 | 161200 | 153.7 | 29.1 | 4.45 |
| zstd | 156100 | 148.8 | 15.9 | 5.62 |
"""
from __future__ import annotations
import json
import os
import resource
import sys
import time
from confluent_kafka import Producer
BOOTSTRAP = sys.argv[2] if len(sys.argv) > 2 else "127.0.0.1:9092"
N = int(sys.argv[1]) if len(sys.argv) > 1 else 100_000
TOPIC = "learn.04.bench"
CODECS = ["none", "gzip", "snappy", "lz4", "zstd"]
def make_payload(i: int) -> bytes:
"""约 1KB 的 JSON 风格 payload,重复字段多以放大压缩比差异。"""
return json.dumps(
{
"event_id": f"evt-{i:08d}",
"user_id": i % 1000,
"action": "purchase",
"category": "electronics",
"currency": "CNY",
"tags": ["promotion", "vip", "newuser", "marketing"] * 5,
"description": "the quick brown fox jumps over the lazy dog. " * 12,
"ts": int(time.time() * 1000),
}
).encode()
def get_cpu_time() -> float:
r = resource.getrusage(resource.RUSAGE_SELF)
return r.ru_utime + r.ru_stime
def bench(codec: str) -> dict:
producer = Producer(
{
"bootstrap.servers": BOOTSTRAP,
"client.id": f"ch4-bench-{codec}",
"acks": "all",
"enable.idempotence": True,
"compression.type": codec,
"linger.ms": 20,
"batch.size": 1048576,
"queue.buffering.max.messages": 500_000,
"queue.buffering.max.kbytes": 1_048_576,
}
)
print(f"\n--- codec = {codec} ---")
cpu0 = get_cpu_time()
t0 = time.perf_counter()
sent = 0
failed = 0
def cb(err, msg):
nonlocal failed
if err:
failed += 1
for i in range(N):
while True:
try:
producer.produce(TOPIC, value=make_payload(i), key=str(i % 50).encode(), on_delivery=cb)
break
except BufferError:
producer.poll(0.5)
sent += 1
if i % 5000 == 0:
producer.poll(0)
remain = producer.flush(120)
elapsed = time.perf_counter() - t0
cpu_used = get_cpu_time() - cpu0
msg_size_avg = sum(len(make_payload(i)) for i in range(min(100, N))) / min(100, N)
bytes_total = sent * msg_size_avg
return {
"codec": codec,
"sent": sent - remain,
"failed": failed,
"elapsed": elapsed,
"msg_per_s": (sent - remain) / elapsed,
"mb_per_s": bytes_total / elapsed / 1024 / 1024,
"cpu_s": cpu_used,
}
def main() -> None:
print(f"== Compression benchmark ==")
print(f"N = {N}, msg ~1KB each, acks=all + idempotence + linger=20")
results = []
for c in CODECS:
r = bench(c)
results.append(r)
print(f" -> sent={r['sent']} failed={r['failed']} {r['elapsed']:.2f}s")
print("\n=== summary ===")
print(f"{'codec':<8} {'msgs/s':>10} {'MB/s':>8} {'cpu_s':>8} {'failed':>8}")
for r in results:
print(f"{r['codec']:<8} {r['msg_per_s']:>10.0f} {r['mb_per_s']:>8.1f} {r['cpu_s']:>8.2f} {r['failed']:>8d}")
print("\nTip: 实际数据中重复模式越多,gzip / zstd 的压缩比越大;"
"随机 / 加密数据用 none 或 lz4 即可,硬压缩纯浪费 CPU。")
if __name__ == "__main__":
main()python
"""
第 4 章 - 幂等 Producer 演示
=================================
演示要点:
1) 开启 enable.idempotence=true 后,Producer 启动会先发
InitProducerIdRequest 拿到 PID + Epoch(看 librdkafka 日志)。
2) 即使同一条消息被 client 在重试时发出多次,broker 也会按
(PID, partition, sequence) 做去重,最终 log 里只出现一次。
3) 开了幂等会自动强制:acks=all、max.in.flight <= 5、retries > 0。
下面的代码故意把 acks=1 注释掉,演示「即使你写了 acks=1 也会被改回 all」。
可以用 kafka-dump-log.sh 验证:
kafka-dump-log.sh \\
--files /tmp/kraft-combined-logs/learn.04.idempotent-0/00000000000000000000.log \\
--print-data-log
看每条消息都有 producerId / sequence 字段。
运行:
bash ../init.sh
python idempotent_producer.py
"""
from __future__ import annotations
import sys
import time
from confluent_kafka import Producer
BOOTSTRAP = sys.argv[1] if len(sys.argv) > 1 else "127.0.0.1:9092"
TOPIC = "learn.04.idempotent"
N = 50
def on_delivery(err, msg):
if err is not None:
print(f" FAIL: {err}")
else:
print(
f" OK : P{msg.partition()}@{msg.offset()} "
f"key={msg.key().decode()} val={msg.value().decode()}"
)
def main() -> None:
print("[ idempotent producer ] enable.idempotence=true")
print("librdkafka will enforce: acks=all, max.in.flight<=5, retries>0\n")
p = Producer(
{
"bootstrap.servers": BOOTSTRAP,
"client.id": "ch4-idempotent-producer",
"enable.idempotence": True,
"linger.ms": 5,
"compression.type": "zstd",
}
)
print(f"Sending {N} messages to {TOPIC}")
for i in range(N):
p.produce(
TOPIC,
key=f"order-{i % 5}".encode(),
value=f"payload-#{i}".encode(),
on_delivery=on_delivery,
)
p.poll(0)
time.sleep(0.01)
p.flush(15)
print("\nDone. Inspect the log with kafka-dump-log.sh; each record carries producerId / sequence.")
if __name__ == "__main__":
main()python
"""
第 4 章 - 同步发送演示
=================================
confluent-kafka 没有原生的同步 send(),本文件演示
如何用 threading.Event + delivery callback 包装同步语义。
适用场景:
- 关键消息必须 ack 后才允许业务继续(订单写完才扣款)
- 单元测试需要确定性
性能:吞吐降到 1k-5k msg/s,比异步低 1-2 个数量级。
运行:
bash ../init.sh
python sync_producer.py # 默认发 100 条
python sync_producer.py 500 myhost:9092
"""
from __future__ import annotations
import sys
import threading
import time
from confluent_kafka import Producer
BOOTSTRAP = sys.argv[2] if len(sys.argv) > 2 else "127.0.0.1:9092"
N = int(sys.argv[1]) if len(sys.argv) > 1 else 100
TOPIC = "learn.04.demo"
class SyncSender:
"""把异步 Producer 包成「单条同步」。"""
def __init__(self, producer: Producer):
self._p = producer
def send(self, topic: str, value: bytes, key: bytes | None = None, timeout: float = 10.0):
done = threading.Event()
result: dict = {}
def cb(err, msg):
result["err"] = err
result["msg"] = msg
done.set()
self._p.produce(topic, key=key, value=value, on_delivery=cb)
# 主动驱动:让 IO 线程把 batch 发出去
self._p.poll(0)
if not done.wait(timeout):
raise TimeoutError(f"send timeout after {timeout}s")
if result["err"] is not None:
raise RuntimeError(f"delivery failed: {result['err']}")
return result["msg"]
def main() -> None:
producer = Producer(
{
"bootstrap.servers": BOOTSTRAP,
"client.id": "ch4-sync-producer",
"acks": "all",
"enable.idempotence": True,
# 关键:同步语义下不要再攒批,linger=0 让消息立即发出
"linger.ms": 0,
}
)
sender = SyncSender(producer)
print(f"🚀 同步发送 {N} 条到 {TOPIC}")
start = time.perf_counter()
for i in range(N):
msg = sender.send(
TOPIC,
value=f"sync-msg-{i}".encode(),
key=f"k{i % 5}".encode(),
)
if i < 5 or i == N - 1:
print(
f" [{i:>4}] partition={msg.partition()} offset={msg.offset()} "
f"latency_ms={(time.perf_counter() - start) * 1000 / (i + 1):.2f}"
)
elapsed = time.perf_counter() - start
print(
f"\n📊 完成 {N} 条;总耗时 {elapsed:.2f}s;"
f"吞吐 {N / elapsed:.0f} msg/s;平均单条 {elapsed / N * 1000:.2f} ms"
)
producer.flush(5)
if __name__ == "__main__":
main()async_producer.py ↗ · benchmark_compression.py ↗ · idempotent_producer.py ↗ · sync_producer.py ↗