阶段一 · 基础与核心模型

消息可靠性:不丢,不重,顺序性

一句话总结

消息可靠性有三个层次:不丢(At-least-once)、不重(Exactly-once)、有序.它们不是 Kafka 单方面保证的,而是 Producer 配置 + Broker 配置 + Consumer 处理逻辑三方配合的结果.
本篇把端到端链路拆开,逐段分析哪里会丢、哪里会重、怎么保证顺序.

消息传递语义:三种级别

语义 含义 代价
At-most-once 最多投递一次,可能丢,不会重 最简单,性能最高
At-least-once 至少投递一次,不会丢,可能重 需要重试 + Consumer 幂等
Exactly-once 恰好一次,不丢不重 最复杂,需要事务或幂等机制

Kafka 默认配置(franz-go)实现的是 At-least-once:Producer 幂等保证不丢,但 Consumer 侧可能重复处理.要做到端到端 Exactly-once,需要额外机制.


第一段:Producer → Broker(不丢)

会丢的场景

场景 原因 解法
acks=0 发出去不等回复,网络丢包就没了 设 acks=all
acks=1 + Leader 宕机 Leader 写入后还没同步到 Follower 就挂了 设 acks=all + min.insync.replicas≥2
Producer 内存缓冲区断电 消息还没发出去,进程挂了 Outbox 模式(上游先落库)
网络超时 + 不重试 消息发了但 ack 丢失,Producer 认为失败却不重试 开启重试(franz-go 默认无限重试可恢复错误)

不丢的标准配置

client, _ := kgo.NewClient(
kgo.SeedBrokers("localhost:9092"),
kgo.RequiredAcks(kgo.AllISRAcks()), // acks=all
// 幂等默认开启 → 自动重试 + 去重
// min.insync.replicas 在 Broker 端配置(server.properties)
)

Broker 端:

# server.properties
min.insync.replicas=2
default.replication.factor=3

三者配合的效果:

Producer(acks=all) → Broker(3 副本,至少 2 个 ISR 确认才返回 ack)
↓
最多容忍 1 个副本故障,消息不丢
如果 2 个副本同时故障 → Broker 拒绝写入(NotEnoughReplicas)
→ 宁可不可用,也不丢数据

第二段:Broker 存储(不丢)

消息写入 Broker 后是否安全?取决于:

风险 原因 Kafka 的保护
单机磁盘损坏 硬件故障 副本机制:数据在多个 Broker 上有拷贝
所有副本同时丢失 机房级故障 跨机架部署(broker.rack配置) + 异地多集群(MirrorMaker)
消息过期被删 retention 到期 调大 retention 或用 compact 策略

只要副本因子 ≥ 2 且副本分布在不同故障域,Broker 端的数据丢失概率极低.


第三段:Broker → Consumer(不丢 + 可能重复)

这是最容易出问题的一段.Consumer 拉到消息后,有一个处理窗口:

PollFetches() → 拿到 [msg-A, msg-B, msg-C]
│
▼ 处理 msg-A ✓
▼ 处理 msg-B ✓
▼ 处理 msg-C ... 进程 crash!
│
└── offset 还没提交!

重启后:从上次提交的 offset 继续 → msg-A, msg-B, msg-C 全部重新消费

不丢:先处理后提交

// 正确:处理完再提交
fetches := client.PollFetches(ctx)
processAll(fetches) // 全部处理完
client.CommitUncommittedOffsets(ctx) // 再提交

// 错误:先提交再处理(At-most-once,可能丢)
fetches := client.PollFetches(ctx)
client.CommitUncommittedOffsets(ctx) // 先推进 offset
processAll(fetches) // 处理过程中 crash → 消息丢了

重复的根源

"先处理后提交"保证了不丢,但带来了重复:

处理 msg-A ✓ → 处理 msg-B ✓ → 提交 offset ... 提交失败(网络抖动)!
重启后:从上次成功的 offset 继续 → msg-A, msg-B 被重新处理

这就是 At-least-once 的本质:不丢和不重在"先处理后提交"模型下不可兼得,除非引入额外机制.


消费端去重:幂等处理

既然 At-least-once 下消息可能重复到达,Consumer 侧需要幂等(idempotent)处理——同一条消息处理多次和处理一次的最终效果一样.

策略一:业务天然幂等

有些操作天然幂等,不需要额外处理:

// "设置用户状态为 VIP" — 执行多次效果一样
db.Exec("UPDATE users SET status = 'vip' WHERE id = ?", userID)

// "覆盖写入 Redis" — 天然幂等
redis.Set(ctx, "user:123:balance", "100.00", 0)

策略二:唯一键去重

// 用消息的唯一标识(topic + partition + offset)作为幂等键
dedupeKey := fmt.Sprintf("%s-%d-%d", r.Topic, r.Partition, r.Offset)

// 方案 A:数据库唯一索引
_, err := db.Exec(
"INSERT INTO processed_events (dedup_key, ...) VALUES (?, ...) ON CONFLICT DO NOTHING",
dedupeKey,
)

// 方案 B:Redis SET NX + 过期时间
ok, _ := redis.SetNX(ctx, "dedup:"+dedupeKey, "1", 24*time.Hour).Result()
if !ok {
return nil // 已处理过,跳过
}

策略三:业务操作与 offset 在同一事务

tx, _ := db.Begin()

// 业务写入
tx.Exec("INSERT INTO orders ...", ...)

// 把 offset 也写入同一个库(自己维护 offset,不用 __consumer_offsets)
tx.Exec(
"INSERT INTO kafka_offsets (group_id, topic, partition, offset) VALUES (?,?,?,?) "+
"ON CONFLICT (group_id, topic, partition) DO UPDATE SET offset = ?",
groupID, r.Topic, r.Partition, r.Offset, r.Offset,
)

tx.Commit()
// 事务保证:业务和 offset 要么一起成功,要么一起回滚

这是真正的 Consumer 端 Exactly-once:即使重复消费,数据库里的 offset 已经推进过了,事务会被忽略.


端到端 Exactly-once:Kafka 事务

Kafka 0.11+ 提供了事务 API,支持跨分区原子写入 + Consumer offset 原子提交:

一个事务内:
1. 从 input-topic 读消息
2. 处理
3. 写结果到 output-topic
4. 提交 input-topic 的 offset
→ 以上 4 步要么全部成功,要么全部回滚

这是 Kafka Streams(Java)的核心能力.在 Go 中,franz-go 支持事务 API:

client, _ := kgo.NewClient(
kgo.SeedBrokers("localhost:9092"),
kgo.TransactionalID("order-processor-tx-0"), // 开启事务
kgo.ConsumerGroup("order-processor"),
kgo.ConsumeTopics("input-orders"),
kgo.DisableAutoCommit(),
)

// 事务循环
for {
fetches := client.PollFetches(ctx)

if err := client.BeginTransaction(); err != nil {
log.Fatal(err)
}

// 处理并生产到下游 topic
fetches.EachRecord(func(r *kgo.Record) {
result := transform(r)
client.Produce(ctx, &kgo.Record{
Topic: "output-orders",
Key: r.Key,
Value: result,
}, nil)
})

// 原子提交:下游消息 + 上游 offset 一起提交
if err := client.EndTransaction(ctx, kgo.TryCommit); err != nil {
client.EndTransaction(ctx, kgo.TryAbort)
continue
}
}

事务的适用场景

Kafka 事务解决的是 consume-transform-produce 模式的 EOS(Exactly-Once Semantics).
如果你的 Consumer 是写外部系统(数据库、Redis、ES),事务帮不了你——因为外部系统不参与 Kafka 事务.此时回退到"幂等处理"策略.


顺序性保证

顺序性在前面篇幅已多次提到,这里做一个完整总结:

Producer 端顺序

条件 顺序保证
同一 Key,幂等 Producer ✓ 分区内严格有序(Broker 按 sequence number 排序)
同一 Key,非幂等 + max.in.flight > 1 ✗ 重试可能乱序
不同 Key ✗ 落入不同 Partition,无跨分区顺序

Consumer 端顺序

条件 顺序保证
单 Consumer 串行处理 ✓ 按 offset 顺序处理
按 Partition 分 goroutine,Partition 内串行 ✓ 每个 Partition 内有序
Worker Pool 打散处理 ✗ 无序(需要业务容忍或自行排序)

端到端有序的完整条件

1. Producer 设置 Key(同业务实体的消息落同一 Partition)
2. 幂等 Producer 开启(防重试乱序)
3. Consumer 在 Partition 粒度串行处理(不在同一 Partition 内并发)

满足以上三点,同一个 Key 的消息从生产到消费严格有序.


总结:可靠性配置速查

目标 Producer 配置 Broker 配置 Consumer 策略
不丢 acks=all, 幂等开启 min.insync.replicas=2, replication.factor=3 先处理后提交
不重 (Producer 幂等保证 Broker 端不重) — 业务幂等 / 唯一键去重 / 事务
有序 设 Key + 幂等 — Partition 内串行处理
端到端 EOS TransactionalID transaction.state.log 配置 事务消费(仅 consume-produce 模式)

下一篇搭建本地开发环境:docker-compose 起 3-Broker 集群 + kafka-ui,常用 CLI 命令,监控指标入门.