一句话总结
Producer 把消息发给 Kafka 的过程不是简单的"发出去就完了"——中间经历了序列化 → 分区路由 → 缓冲批量 → 网络发送 → 等待 ack 五个阶段.
理解这条链路,才能在"吞吐、延迟、可靠性"三角中做出正确取舍.
Producer 发送全链路
一条消息从你的 Go 代码到最终落盘,经历的完整路径:
|
Go 客户端(如 franz-go)内部也遵循这个模型,只是实现细节不同.接下来逐层拆解.
序列化:消息就是 []byte
Kafka Broker 不关心消息内容,它只存 key 和 value 两段字节数组.序列化是客户端的事.
| 方案 | 优点 | 缺点 | Go 生态 |
|---|---|---|---|
| JSON | 人类可读,调试方便 | 体积大,无 schema 约束 | encoding/json |
| Protobuf | 体积小,有 schema 演进 | 需要 .proto 管理 | google.golang.org/protobuf |
| Avro | Schema Registry 原生支持 | Go 库生态弱于 Java | linkedin/goavro |
| 裸字节 | 零开销 | 无结构,仅适合简单场景 | 直接 []byte |
推荐
Go 后端服务间通信已经用 Protobuf 的,Kafka 消息也用 Protobuf 保持一致.
后续 Schema 管理篇会详细讲 Schema Registry + Protobuf 的配合.
分区路由:消息写到哪个 Partition
这是 Producer 最核心的决策点.规则:
|
各策略对比
| 策略 | 触发条件 | 效果 | 适用场景 |
|---|---|---|---|
| 指定 Partition | 代码显式设置 | 完全控制 | 极少用,测试或特殊路由 |
| Key Hash | 设置了 Key | 相同 Key 永远落同一 Partition | 需要顺序保证:同用户、同订单 |
| Sticky(默认) | Key 为 nil | 攒满一批后随机换 Partition | 无序要求,追求均匀分布和批量效率 |
| Round Robin | 老版本默认 | 逐条轮转 | 已被 Sticky 取代,批量效率差 |
| 自定义 | 实现 Partitioner 接口 | 任意逻辑 | 地域路由、优先级分区 |
Key Hash 的细节
默认哈希算法是 murmur2(Kafka 官方客户端统一使用,保证跨语言一致):
|
分区数变更会打乱路由
hash(key) % N 依赖分区数 N.如果你从 6 个分区扩到 12 个,同一个 Key 会路由到不同 Partition,原有的顺序保证被打破.
所以:对顺序敏感的 Topic,创建时就要规划好分区数,尽量不扩.
Sticky Partitioner:为什么替代了 Round Robin
旧版 Round Robin 对无 Key 消息逐条轮转,导致每条消息落不同 Partition,batch 永远凑不满:
|
Go 客户端(franz-go)默认就是 Sticky 行为,无需额外配置.
缓冲与批量:RecordAccumulator
Producer 不会每条消息都发一次网络请求——它在内存里按 Partition 分桶缓冲,凑够一批再发.
|
两个关键参数
| 参数 | 含义 | 默认值 | 调优方向 |
|---|---|---|---|
batch.size |
单个 batch 的字节上限 | 16 KB | 调大 → 吞吐高,延迟略高 |
linger.ms |
batch 没满时最多等多久就发 | 0 ms | 调大(5-100ms)→ 有机会凑更大批 |
linger.ms = 0 表示"有消息就立刻发",此时 batch.size 几乎不起作用(除非短时间内涌入大量消息恰好填满).
类比 Go 的 bufio.Writer:你 Write() 多次,它在内存攒到缓冲区满了才真正调 syscall.
batch.size 是缓冲区大小,linger.ms 是"最多等多久就 Flush 一次".
buffer.memory:全局内存上限
所有 Partition 的 batch 共享一个内存池,默认 32 MB.如果 Producer 发送速率远超 Broker 消费速率,内存池耗尽时:
- Java 客户端:阻塞,等待
max.block.ms后抛异常 - franz-go:
Produce()返回 error 或阻塞(取决于配置)
这是 Producer 端的背压机制——防止无限堆积撑爆内存.
acks:可靠性的旋钮
消息发到 Broker 后,Producer 需要等一个 ack 回复才算"发送成功".acks 参数控制等到什么程度:
| acks | 含义 | 持久性 | 延迟 | 吞吐 |
|---|---|---|---|---|
0 |
不等回复,发出去就算成功 | 可能丢(网络丢包就没了) | 最低 | 最高 |
1 |
Leader 写入本地日志即回复 | Leader 挂了可能丢(follower 还没同步) | 中等 | 中等 |
all(-1) |
Leader 等所有 ISR 副本确认 | 最强(ISR 全挂才可能丢) | 最高 | 最低 |
acks=all 的真实含义
"all"不是所有副本,是所有 ISR(In-Sync Replicas) 成员.ISR 是与 Leader 保持同步的副本子集.配合 min.insync.replicas 使用:
|
生产环境推荐配置
|
这组配置意味着:最多容忍 1 个副本故障,同时保证数据不丢.
是可靠性与可用性的最佳平衡点.
重试与幂等
网络不可靠,发送失败怎么办?
自动重试
| 参数 | 含义 | 默认值 |
|---|---|---|
retries |
最大重试次数 | Java 客户端:MAX_INT;franz-go:无限重试可恢复错误 |
retry.backoff.ms |
重试间隔 | 100 ms |
delivery.timeout.ms |
从调用 produce 到最终放弃的总时限 | 120000 ms(2 分钟) |
重试带来的问题:消息重复与乱序
|
乱序场景:
|
解决方案:幂等 Producer
开启幂等后,Producer 给每条消息附加 <ProducerID, SequenceNumber>,Broker 据此去重和排序:
|
| 配置 | 效果 |
|---|---|
enable.idempotence = true |
Broker 按 PID+Seq 去重,保证分区内 exactly-once 写入 |
自动设置 acks = all |
幂等要求 acks=all |
自动设置 max.in.flight = 5 |
Broker 端按 Seq 排序,即使乱序到达也能还原顺序 |
franz-go 的默认行为
franz-go 默认开启幂等、acks=all、自动重试.
也就是说:用 franz-go 开箱即用,不需要你操心去重和乱序问题.
这是选它而非 sarama 的一个重要理由.
Go 代码:Producer 最小示例
用 franz-go 展示核心 API,后续实战篇会展开完整工程结构:
|
几个关键点:
| 代码 | 说明 |
|---|---|
kgo.SeedBrokers |
初始连接地址,客户端会自动发现完整集群 |
Key: []byte("user-123") |
设置 Key → hash 路由,同用户消息落同一 Partition |
ProduceSync |
阻塞等待结果;高吞吐场景用异步 Produce + callback |
client.Close() |
刷空缓冲区中未发送的 batch,等待 ack |
异步 vs 同步发送
| 模式 | API | 吞吐 | 延迟感知 | 适用场景 |
|---|---|---|---|---|
| 同步 | ProduceSync() |
低(每条等 ack) | 立即知道成功失败 | 低频但关键的事件 |
| 异步 + callback | Produce(ctx, record, callback) |
高(不阻塞调用方) | callback 中处理 | 高吞吐主流用法 |
| Fire-and-forget | Produce(ctx, record, nil) |
最高 | 不关心结果 | 日志、指标等允许丢的场景 |
异步发送的典型写法:
|
小结:Producer 关键参数速查
| 参数 | 作用 | 推荐值(生产环境) |
|---|---|---|
acks |
可靠性等级 | all |
batch.size |
单批大小 | 64KB-256KB(视消息大小) |
linger.ms |
凑批等待时间 | 5-20ms |
buffer.memory |
全局缓冲上限 | 64MB-128MB |
compression.type |
压缩算法 | lz4(低 CPU)或 zstd(高压缩比) |
enable.idempotence |
幂等去重 | true(franz-go 默认) |
max.in.flight.requests |
并发在途请求数 | 5(幂等模式下的上限) |
下一篇进入 Consumer 端:Consumer Group 的协作机制、rebalance 原理、offset 提交策略.