Skip to content

第 4 章 Producer 深入:把一条 send() 拆成三十步

目标读者:第 3 章已经会用 kafka-console-producer.shconfluent-kafkaproduce() 发消息,但被问到「acks=all 到底等几个副本」「linger.ms=20batch.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=5ms168,0001835%
异步 + linger=0ms92,000950%
同步(每条 wait)1,2000.88%

结论:没有特殊业务原因就用异步 + 回调,再用 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 三个参数各自含义

参数默认值含义
retriesInteger.MAX_VALUE(3.0+)单个 batch 的最大重试次数;只决定「最多试几次」
retry.backoff.ms100两次重试之间的最小间隔
delivery.timeout.ms120000从 produce() 入队到 ack(含全部重试)的总超时
request.timeout.ms30000单次 ProduceRequest 等响应的超时
max.in.flight.requests.per.connection5单个 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=allretries>0,不需要你手动配。


5. linger.ms / batch.size / compression:吞吐三剑客

5.1 它们的关系

linger.msbatch.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.msbatch.size吞吐 (MB/s)平均延迟 (ms)p99 延迟 (ms)网络包数
01638422512极多(每条一包)
51638489818
20163841542340
10016384168105130极少
565536132922
51MB1541125

观察结论:

  1. linger.ms=0 即「来一条发一条」,吞吐惨不忍睹但延迟最小。
  2. linger.ms 在 5-20 之间是大多数业务的甜蜜点:吞吐 6-7 倍,延迟只多个十几毫秒。
  3. batch.size 单方面调大,效果不如 linger.ms 明显(因为没消息就攒不起来)。
  4. 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
none1.0×15415428%
gzip6.8×921488%
snappy3.1×1384541%
lz43.0×1625435%
zstd5.5×1562852%

经验法则

  • 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=none91,0008默认配置
改 acks=all (RF=3, min.isr=2)38,00022可靠性升一级,吞吐降一半
加 linger.ms=20142,00035攒批起作用
加 compression=zstd156,00038网络省 5x
开 enable.idempotence152,00040吞吐略降,去重保障
优化后:acks=all + linger=20 + zstd + idempotence152,00040吞吐 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 后:

  1. Producer 启动时,向某个 Broker 发 InitProducerIdRequest,broker 给它分配一个全局唯一的 ProducerId(PID)+ Epoch
  2. 每个 (PID, partition) 维护一个自增序列号 sequence number,每条消息从 0 递增。
  3. 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=allmax.in.flight ≤ 5retries > 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_AVAILABLENOT_LEADER_FOR_PARTITIONNETWORK_EXCEPTIONREQUEST_TIMED_OUT自动重试(受 retries / delivery.timeout 限制)
不可重试RECORD_TOO_LARGEINVALID_TOPICUNKNOWN_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 ProducerRabbitMQ ProducerRocketMQ Producer
写入模型批量异步(每分区 batch + Sender 线程)单条 publish + confirm同步 / 异步 / OneWay 三模式
顺序保证单分区严格有序(开幂等才保险)单队列有序MessageQueueSelector + 同步发送可保证
Exactly Once幂等 (PID+seq) + 事务 (transactional.id)RabbitMQ 5+ 才有 streams 支持半消息事务
压缩batch 级 (gzip/snappy/lz4/zstd)单条 / 不内置压缩单条 zip
路由按 Key HashExchange + 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.demolearn.04.idempotentlearn.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 各自的语义、典型丢消息场景、和性能影响?

考察点:可靠性核心,最高频题之一。

答案

  1. acks=0:Producer 发出请求后不等任何响应就认为成功。
    • 可能丢:Broker 没收到 / 收到但写失败 / Leader 切换全丢。
    • 吞吐最高,但重试无意义(没响应就没失败信号)。
    • 适用:可丢的指标采样、A/B 实验事件。
  2. acks=1:Leader 写入本地 log 就回 ack。
    • 仍可能丢:Leader 写完但 Follower 还没拉到 → Leader 宕机 → 新 Leader 没这条消息。
    • 适用:内部日志、能容忍极少量丢失。
  3. acks=-1 / acks=all:等 ISR 中所有副本都同步完才 ack。
    • 不丢的前提min.insync.replicas ≥ 2,否则 ISR 只剩 Leader 时退化为 acks=1。
    • 业界推荐组合:RF=3 + min.insync.replicas=2 + acks=all
    • 适用:金融、订单、事务、Exactly Once。
  4. 吞吐对比(同环境):acks=0 ≈ 250k msg/s,acks=1 ≈ 170k,acks=all ≈ 90k。
  5. 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 一定会乱序吗?开了幂等还会吗?

考察点:重试 / 顺序 / 幂等的交互,深度题。

答案

  1. 不开幂等max.in.flight > 1 时,多个请求并发在网络上,一旦前一个失败重试,可能比后一个晚到 → 同一分区内消息乱序
    • 严格保序的传统做法:max.in.flight=1,吞吐降到 1/5。
  2. 开幂等 (enable.idempotence=true):
    • Broker 维护 (PID, partition) 维度的 last_seq,每条消息带自增 sequence number。
    • 即使 client 把 5 个请求并发发出,broker 也会按 seq 顺序检查:seq != last+1 时拒绝。
    • 失败重试时,broker 看到老 seq 直接丢弃(去重),不影响顺序。
    • client 端会自动把 max.in.flight 限制在 ≤ 5(库内部强制,超过会启动失败)。
  3. 结论:开幂等 + max.in.flight ≤ 5 是「既保序又高吞吐」的最优组合,不开幂等的话只能在「max.in.flight=1(保序)」和「>1(高吞吐但乱序)」之间二选一。
  4. 跨分区:同一 Producer 写多个分区时,Kafka 不保证跨分区顺序——这是分区模型的本质,不可改变。需要全局顺序就只能用单分区(牺牲并行度)。

加分项:提到事务 Producer(transactional.id)让多个分区的写入原子化,但顺序仍是分区内的;Kafka Streams 内部就是用「按 Key 路由保证 partition 内顺序」实现一致性。


Q3:linger.ms / batch.size 是「与」还是「或」?怎么调?

考察点:吞吐调优基本功。

答案

  1. 「或」关系:任意一个先满足就触发发送。
    • batch 当前大小 ≥ batch.size → 发车。
    • linger.ms 计时器超时 → 发车。
  2. linger.ms 的取舍
    • 0:来一条发一条,吞吐惨(每条消息都付出固定开销),延迟最低。
    • 5-20:大多数业务的甜蜜点,吞吐能提 5-10x,延迟只多十几毫秒。
    • >50:吞吐基本到顶,再加只是多攒延迟。
  3. batch.size 的取舍
    • 默认 16384(16KB)。
    • 调大(64KB-1MB)能让单批包含更多消息,配合 linger 更高效;超过 1MB 通常没收益(被 broker 端 message.max.bytes 限制)。
  4. 「调多少」实操
    • 业务延迟容忍 50ms 以内 → linger=10,batch.size=64KB,吞吐已经能跑满千兆网。
    • 能容忍 200ms+ 的离线管道 → linger=100,batch.size=1MB,吞吐到顶。
  5. 配套:开 compression=zstd 才能让大 batch 的网络收益最大化(小 batch 压缩比低)。
  6. 指标:盯 batch-size-avg,如果远小于 batch.size,说明流量太稀疏,调大 batch 没用,应该调大 linger.ms 或合并业务。

加分项:提到 librdkafka 还有 queue.buffering.max.ms 是 linger.ms 的别名;Java 中 KafkaProducermetrics() 能拿到 record-queue-time-avg 直观看到消息在 accumulator 里待了多久。


Q4:enable.idempotence=true 的实现原理?能不能跨 Producer 重启去重?

考察点:幂等的内部机制。

答案

  1. 核心三件套
    • PID(Producer ID):Producer 启动时向 broker 申请,全局唯一。
    • Epoch:每次 PID 被「重新认领」(事务场景下复用 transactional.id 时)就 +1,旧 Epoch 的请求被拒。
    • 序列号 (sequence number):按 (PID, partition) 维度从 0 自增,每条消息带一个。
  2. broker 端去重
    • 每个 (PID, partition) 缓存最近 5 个 seq(受 max.in.flight ≤ 5 限制)。
    • 收到 seq == last_seq + 1 才接受,seq <= last_seq 视为重复直接丢弃返回成功,seq > last_seq + 1 视为消息丢失返回 OutOfOrderSequenceException
  3. 解决的问题:网络抖动导致 client 重试 → broker 去重 → 单 Session 内不重不丢
  4. 跨 Session 不行
    • Producer 重启会重新 InitProducerIdRequest,拿到新 PID
    • broker 视为完全不同的 Producer,去重历史失效。
    • 想跨重启幂等:用事务 Producer,配 transactional.id,broker 会让同一 transactional.id 复用 PID 并自增 Epoch,把旧 Epoch 的「僵尸消息」屏蔽。
  5. 强制约束:开了 idempotence 会自动强制
    • acks=all
    • max.in.flight.requests.per.connection ≤ 5
    • retries > 0
    • 违反会在客户端启动时报 ConfigException
  6. 代价:吞吐降低 < 5%,broker 内存多缓存几个 seq,几乎免费,3.0+ 强烈建议默认开。

加分项:提到 Kafka 3.0 起 enable.idempotence=true 是默认值;__transaction_state 内部 topic 存储事务状态,去重的 PID/Epoch 信息存在 broker 内存(重启会从 __transaction_state 恢复事务上下文,但单纯幂等的 PID 不持久化)。


Q5:delivery.timeout.msrequest.timeout.msretriesretry.backoff.ms 怎么配合?

考察点:超时 + 重试模型。

答案

  1. 各自含义
    • 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),是兜底总闸门。
  2. 关系约束
    delivery.timeout.ms ≥ linger.ms + request.timeout.ms × (retries + 1)
                          + retry.backoff.ms × retries
    设小了 delivery.timeout.ms 会让 retries 失效。
  3. 触发失败的条件:先到先触发——
    • 重试次数超过 retries
    • 总耗时超过 delivery.timeout.ms
  4. 3.0 起的最佳实践retries 不用动(默认极大),重点调 delivery.timeout.ms 控制业务能接受的最长等待。
  5. 拓展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?

考察点:内部数据结构、内存控制。

答案

  1. RecordAccumulator 是 Producer 在客户端进程内的内存缓冲区
    • 按 (Topic, Partition) 维度组织,每个分区一个双端队列 Deque<ProducerBatch>
    • 业务线程 produce() 把消息追加到对应分区的最后一个 batch,满了就新建一个 batch。
    • Sender 后台线程从这里 drain(按 broker 维度合并)发出。
  2. 总大小由 buffer.memory(默认 32MB)控制,所有分区的 batch 共享。
  3. BufferError / TimeoutException 的成因:业务发送速度 > Sender 发送速度 + broker 处理速度,buffer 被打满。
  4. 解决思路
    • 加大 bufferbuffer.memory=128MB(注意会占用 client 内存)。
    • 加快发送:调大 max.in.flight、开压缩、加 broker 数量、扩分区。
    • 业务限速:在业务侧加令牌桶,避免突发超过下游能力。
    • 改异步队列:业务线程 try/except BufferError,sleep 后重试,让 Sender 喘息。
  5. librdkafka 的 Python 客户端queue.buffering.max.messages(默认 100k 条)、queue.buffering.max.kbytes(默认 1GB)控制本地队列;produce()BufferError 时调 producer.poll(1) 触发回调释放空间。
  6. 监控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 ↗