阶段一 · 基础与核心模型

生产者模型与分区策略

一句话总结

Producer 把消息发给 Kafka 的过程不是简单的"发出去就完了"——中间经历了序列化 → 分区路由 → 缓冲批量 → 网络发送 → 等待 ack 五个阶段.
理解这条链路,才能在"吞吐、延迟、可靠性"三角中做出正确取舍.

Producer 发送全链路

一条消息从你的 Go 代码到最终落盘,经历的完整路径:

你的 Go 代码
│
│ Produce(topic, key, value)
▼
┌─────────────────────────────────────────────┐
│ Producer 客户端内部 │
│ │
│ 1. Interceptors(可选,埋点/链路追踪) │
│ │ │
│ 2. Serializer(key/value 序列化为 []byte) │
│ │ │
│ 3. Partitioner(决定写哪个 Partition) │
│ │ │
│ 4. RecordAccumulator(按 Partition 缓冲) │
│ │ │
│ 5. Sender 线程(批量发送到 Broker) │
│ │ │
└─────────┼───────────────────────────────────┘
│ TCP 请求(ProduceRequest)
▼
Broker (Leader of target Partition)
│
│ 写入本地日志 + 等待副本同步(取决于 acks)
▼
ProduceResponse(success / error)

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 最核心的决策点.规则:

if 消息指定了 Partition:
直接写入该 Partition
elif 消息有 Key:
partition = hash(key) % numPartitions
else:
使用 Sticky Partitioner(攒一批发同一个分区,下一批换)

各策略对比

策略 触发条件 效果 适用场景
指定 Partition 代码显式设置 完全控制 极少用,测试或特殊路由
Key Hash 设置了 Key 相同 Key 永远落同一 Partition 需要顺序保证:同用户、同订单
Sticky(默认) Key 为 nil 攒满一批后随机换 Partition 无序要求,追求均匀分布和批量效率
Round Robin 老版本默认 逐条轮转 已被 Sticky 取代,批量效率差
自定义 实现 Partitioner 接口 任意逻辑 地域路由、优先级分区

Key Hash 的细节

默认哈希算法是 murmur2(Kafka 官方客户端统一使用,保证跨语言一致):

// 伪代码:franz-go 内部逻辑
func partition(key []byte, numPartitions int) int {
h := murmur2(key)
return int(h&0x7fffffff) % numPartitions // 取正数再取模
}

分区数变更会打乱路由

hash(key) % N 依赖分区数 N.如果你从 6 个分区扩到 12 个,同一个 Key 会路由到不同 Partition,原有的顺序保证被打破.
所以:对顺序敏感的 Topic,创建时就要规划好分区数,尽量不扩.

Sticky Partitioner:为什么替代了 Round Robin

旧版 Round Robin 对无 Key 消息逐条轮转,导致每条消息落不同 Partition,batch 永远凑不满:

Round Robin(旧):
msg1 → P0, msg2 → P1, msg3 → P2, msg4 → P0 ...
每个 Partition 的 batch 都只有 1 条,频繁触发小包发送

Sticky(新):
msg1~msg100 → P0(凑满一个 batch)
msg101~msg200 → P1(换一个 Partition 继续攒)
batch 能攒满,网络效率高

Go 客户端(franz-go)默认就是 Sticky 行为,无需额外配置.


缓冲与批量:RecordAccumulator

Producer 不会每条消息都发一次网络请求——它在内存里按 Partition 分桶缓冲,凑够一批再发.

┌── RecordAccumulator ──────────────────────┐
│ │
│ Partition 0: [msg, msg, msg, ...] batch │
│ Partition 1: [msg, msg] batch │
│ Partition 2: [msg, msg, msg, msg] batch │
│ │
└───────────────────────────────────────────┘
│
│ 触发条件:batch 满了 OR linger.ms 到了
▼
Sender → Broker

两个关键参数

参数 含义 默认值 调优方向
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 使用:

replication.factor = 3        (共 3 个副本)
min.insync.replicas = 2 (至少 2 个 ISR 成员才允许写入)
acks = all (等 ISR 全部确认)

效果:
- 正常时 3 副本,ISR=[0,1,2],3 个都确认才回复
- 一个 follower 掉线,ISR=[0,1],2 个确认即可
- 两个 follower 掉线,ISR=[0] 只剩 1 个 < min.insync.replicas
→ Broker 拒绝写入,返回 NotEnoughReplicasException
→ 宁可不可用,也不丢数据

生产环境推荐配置

acks = all
min.insync.replicas = 2
replication.factor = 3

这组配置意味着:最多容忍 1 个副本故障,同时保证数据不丢.
是可靠性与可用性的最佳平衡点.


重试与幂等

网络不可靠,发送失败怎么办?

自动重试

参数 含义 默认值
retries 最大重试次数 Java 客户端:MAX_INT;franz-go:无限重试可恢复错误
retry.backoff.ms 重试间隔 100 ms
delivery.timeout.ms 从调用 produce 到最终放弃的总时限 120000 ms(2 分钟)

重试带来的问题:消息重复与乱序

Producer → Broker:  发送 batch [A, B]
网络超时,Producer 认为失败
Producer → Broker: 重试发送 [A, B]
第一次其实 Broker 已经收到了
结果:Partition 里有两份 [A, B] → 重复!

乱序场景:

max.in.flight.requests.per.connection = 5 (默认)

batch1 发送 → 失败 → 排队重试
batch2 发送 → 成功
batch1 重试 → 成功
结果:落盘顺序是 batch2, batch1 → 乱序!

解决方案:幂等 Producer

开启幂等后,Producer 给每条消息附加 <ProducerID, SequenceNumber>,Broker 据此去重和排序:

// franz-go 默认就开启幂等(RequiredAcks = AllISRAcks 时)
client, _ := kgo.NewClient(
kgo.SeedBrokers("localhost:9092"),
kgo.RequiredAcks(kgo.AllISRAcks()),
// 幂等自动启用,无需额外配置
)
配置 效果
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,后续实战篇会展开完整工程结构:

package main

import (
"context"
"fmt"
"github.com/twmb/franz-go/pkg/kgo"
)

func main() {
client, err := kgo.NewClient(
kgo.SeedBrokers("localhost:9092"),
kgo.DefaultProduceTopic("order-events"),
// 幂等 + acks=all 默认开启
)
if err != nil {
panic(err)
}
defer client.Close()

// 同步发送:等待 ack 返回
record := &kgo.Record{
Key: []byte("user-123"),
Value: []byte(`{"action":"purchase","amount":99.9}`),
}

results := client.ProduceSync(context.Background(), record)
for _, pr := range results {
if pr.Err != nil {
fmt.Printf("发送失败: %v\n", pr.Err)
} else {
fmt.Printf("发送成功: topic=%s partition=%d offset=%d\n",
pr.Record.Topic, pr.Record.Partition, pr.Record.Offset)
}
}
}

几个关键点:

代码 说明
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) 最高 不关心结果 日志、指标等允许丢的场景

异步发送的典型写法:

client.Produce(ctx, record, func(r *kgo.Record, err error) {
if err != nil {
// 所有重试都失败了才会到这里
log.Error("produce failed", "err", err, "key", string(r.Key))
// 根据业务决定:写入死信队列 / 报警 / 降级
return
}
// 成功:r.Partition, r.Offset 已填充
})

小结: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 提交策略.