Kafka 不是一个「有人推消息给你」的队列。它是一份按 Partition 切开、只能追加的分布式日志。
Producer 决定写进哪个 Partition;Consumer 自己 poll,靠两样东西知道下一轮该读什么:组里分到了哪些 Partition,以及每个 Partition 当前的 offset。组内同一时刻一个 Partition 只会交给一个成员,所以不会出现两个人同时啃同一条消息。但崩溃、rebalance、提交失败,仍然会让已经处理过的消息再被读一遍。
所以重复消费要分两层看:协议层保证的是分区独占,不是业务 exactly-once。后者得靠手动提交 offset,以及业务幂等。
代码可以参考 kafka-example。
它到底是什么
Topic 只是分类名。真正存储、真正能并行的单位是 Partition:一段有序、不可变、只追加的日志,每条消息有一个单调递增的 offset,像数组下标。
分区 0
[消息] [消息] [消息] [消息] [消息] [消息]
0 1 2 3 4 5 ← offset
一个 Topic 切成多个 Partition,读写就可以摊到不同 Broker 上。每个 Partition 再复制成若干副本:一个 leader,其余 follower。Producer 只写 leader;Consumer 默认也从 leader 拉(2.4 之后可以按机架从 follower 读,读到的仍不超过该副本的 HW)。follower 负责把日志追上,leader 挂了才有资格顶上。
flowchart LR P[Producer] --> L[Broker 2<br/>P0 leader] L --> C[Consumer] F1[Broker 1<br/>P0 follower] -.->|fetch| L F2[Broker 3<br/>P0 follower] -.->|fetch| L CTRL[Controller / KRaft] --- L CTRL --- F1 CTRL --- F2
Kafka 4.0 已经去掉 ZooKeeper,元数据由 KRaft 法定人数用 Raft 维护。这只影响集群元数据,分区日志本身仍然不是 Raft:还是 leader + ISR + follower pull。后文会回到这一点。
关系和基数可以先记在这张图里:

几个容易混的数量关系:
- 一个 Topic 有 0 个或多个 Producer、0 个或多个 Consumer
- 一个 Topic 被切成 1 个或多个 Partition
- 一个 Partition 在同一个 Consumer Group 内同一时刻只有 1 个 Consumer
- 一个 Consumer 可以同时拉多个 Partition
- 一个 Partition 有 1 个 leader、0 个或多个 follower;每个副本落在一台 Broker 上
消息队列常被拿来做三件事:削峰、解耦、让下游按自己的节奏重放。Kafka 的特别之处是:日志不因消费而消失,直到 retention 或 compaction 把它收走。不同业务用不同 group.id 订阅同一个 Topic,各记各的 offset,互不影响。订单 Topic 同时喂给风控、报表、通知,靠的就是这个,不是什么额外的广播模式。
生产者怎么写进去
一次写入大致是四步。
- 应用构造一条
ProducerRecord:topic、key、value,再加上 timestamp、headers。 - 分区器选定目标 Partition。有 key 时通常按 hash(Java 客户端是 murmur2)路由,同一个 key 会进同一个 Partition,分区内有序。无 key 时,2.4 之后默认是 sticky partitioner:先粘在一个分区上把 batch 攒满,再换下一个,吞吐比纯轮询好。
- 消息发给该 Partition 的 leader。
- Broker 按
acks决定何时返回确认。
acks |
含义 | 风险 |
|---|---|---|
0 |
发出去就当成功 | 吞吐最高,丢了也不知道 |
1 |
leader 写入自己的日志就确认 | leader 在 follower 追上之前挂了,消息会丢 |
all / -1 |
等到当前 ISR 全部写入 | 可靠,延迟更高 |
acks=all 等的不是「集群里每一个副本」,是 当前 ISR 里的每一个。真正托底的是 min.insync.replicas:ISR 人数掉到这个值以下,acks=all 的写入会被直接拒绝。它默认是 1,意味着 follower 全掉光之后,只剩 leader 也能写成功——这一刻可靠性和 acks=1 没区别。生产上常见组合是 replication.factor=3、min.insync.replicas=2,再配上 unclean.leader.election.enable=false,这样 ISR 只剩 leader 时宁可不写,也不假装成功。
重试会引入另一类问题。网络抖一下,Broker 其实写进去了,Producer 没收到 ack,再发一次,日志里就多了一条内容相同、offset 不同的消息。从 0.11 起可以打开幂等生产者:Broker 用 PID + epoch + sequence 认出「这是同一条重试」,把重复的丢掉。Kafka 3.0 之后,Java 客户端默认 enable.idempotence=true,同时默认 acks=all。
幂等生产者只管同一次 Producer 运行、同一个分区上的重试。进程重启会换新的 PID,Broker 认不出这是上一轮的重试。换一个进程再写一遍同样的业务事件,或者跨多个分区要原子提交,它也不管。跨重启的 fencing、以及多分区原子写,要上事务(transactional.id)。
Sarama 里对应的是:
config := sarama.NewConfig()
config.Producer.Idempotent = true
config.Producer.RequiredAcks = sarama.WaitForAll
config.Producer.Retry.Max = 5
config.Net.MaxOpenRequests = 1
MaxOpenRequests 必须压到 1(Sarama 对幂等的约束比 Java 客户端更严,Java 允许 max.in.flight.requests.per.connection <= 5)。否则乱序重试会把 sequence 打乱。
消费者怎么知道该处理什么
Consumer 是 pull,不是 push。它加入某个 Consumer Group,Group Coordinator 按订阅关系和分配策略,把 Partition 派给组员。之后 Consumer 循环 poll(),从自己名下的 Partition、从当前 offset 往后拉一批。
所谓「感知到有活」,不是 Broker 敲它的门,是它自己拿着「分区分配 + offset」去日志里读。
offset 记的是下一条该读的位置。第一次进来、还没有提交记录时,看 auto.offset.reset / Sarama 的 Offsets.Initial:
config.Consumer.Offsets.Initial = sarama.OffsetOldest
// 从分区里还留着的最早消息开始
[0][1][2][3][4][5]
↑
config.Consumer.Offsets.Initial = sarama.OffsetNewest
// 只追新消息,历史全部跳过
[0][1][2][3][4][5]
↑
提交成功的进度写在内部 Topic __consumer_offsets 里,key 是 (group, topic, partition)。每个组、每个分区各记一份。组挂了再拉起来,就从这份进度继续。
默认的自动提交很危险:它按时间间隔提交的是最近一次 poll 到的位置,不是「业务已经处理完」的位置。处理到一半进程没了,offset 却已经越过这些消息,等于丢。反过来,处理完了还没轮到提交就崩了,等于重放。
所以业务消费应该关掉自动提交,处理成功后再标进度:
config.Consumer.Offsets.AutoCommit.Enable = false
session.MarkMessage(msg, "")
Sarama 的 MarkMessage 只是在这次 session 里记下「这条可以提交了」,不是 ack 瞬间落盘。自动提交打开时,后台会按间隔以及 session 结束时把已 Mark 的位移刷出去;关掉之后必须自己调 session.Commit(),否则 rebalance 或进程退出都会把这些 Mark 丢掉,等于没提交。
消费循环本身长这样:
flowchart TD
A[加入 Consumer Group] --> B[拿到 Partition 分配]
B --> C[poll 一批消息]
C --> D{处理成功?}
D -->|是| E[Mark / 手动提交 offset]
E --> C
D -->|否| F[不提交,返回或重试]
F --> C
B -.->|rebalance / 崩溃| G[旧 session 结束]
G --> A
Sarama 把一次分区归属收成 ConsumerGroupSession,三次回调对应一次会话的寿命:
type ConsumerGroupHandler struct{}
func (h *ConsumerGroupHandler) Setup(session sarama.ConsumerGroupSession) error {
// 新分配生效:可以在这里记录自己拿到了哪些分区
return nil
}
func (h *ConsumerGroupHandler) ConsumeClaim(
session sarama.ConsumerGroupSession,
claim sarama.ConsumerGroupClaim,
) error {
for {
select {
case msg, ok := <-claim.Messages():
if !ok {
return nil
}
if err := process(msg); err != nil {
// 不 Mark、不 Commit。返回后这次归属结束,下次从旧 offset 重读
return err
}
session.MarkMessage(msg, "")
session.Commit() // AutoCommit=false 时必须自己提交;生产里可以每 N 条再 Commit
case <-session.Context().Done():
return nil
}
}
}
func (h *ConsumerGroupHandler) Cleanup(session sarama.ConsumerGroupSession) error {
// 分区被收回之前的清理。不要在这里假设还能继续消费
return nil
}
注意两件实现细节:
ConsumeClaim必须把claim.Messages()读干,或者立刻退出让 session 结束。卡在中间不读也不返回,Sarama 的心跳仍可能在后台续着,但分区收不回去、rebalance 也会卡住。- 处理失败就不要
Mark,也不要Commit。Kafka 没有 NACK,不提交就是「下次还从这里开始」。offset 是分区上的一个数字,后面的消息一旦被提交,失败的那条就被跳过了。
组内为什么看起来不重复
多个进程共用同一个 group.id,就是同一个 Consumer Group。对同一个 Topic 的同一个 Partition,组内同一时刻只有一个成员拥有消费权。这是 Kafka 避免组内并发重复的全部手段。

六个 Partition、三个 Consumer,理想情况是每人两个。再加第四个、第五个、第六个还能继续摊薄;第七个开始就会有人空转——并行度的上限是 Partition 数,不是 Consumer 数。
图里写在 Consumer 身上的「Offset: 42」是简化。真实进度是每个 (group, topic, partition) 一份。一个人吃两个分区,就有两条 offset。
再强调一次反面:
- 同一个 Partition 不能被同组两个 Consumer 同时拉
- 一个 Consumer 可以同时拉多个 Partition
- 不同组会各自完整地消费同一份 Topic,这是特性,不是 bug
- 组与组之间的 offset 完全独立
顺序同样只在分区里成立。要让同一个订单的事件保序,就把 order_id 当 key,让它们进同一个 Partition。跨分区没有全局顺序。
Rebalance:重复处理的高发区
组成员变了,分区就要重分。常见触发:
- 有人加入
- 有人主动离开
- 有人崩溃或心跳超时
- Topic 的分区数变了
- 订阅关系变了
旧协议是 stop-the-world:所有人先停消费,交还分区,再一次性重新领。新成员加入的那几秒到几十秒,整组都不干活。较新的客户端可以改用增量协作(Java 里是 CooperativeStickyAssignor):只停那些真正要易手的分区,其余人继续拉。

rebalance 能保证分区最终重新有主,不保证「正在处理的那条」只会被做一次。
| 场景 | 为什么会重复 | 怎么收 |
|---|---|---|
| Consumer 崩溃 | 业务写库成功,offset 还没提交 | 关自动提交;成功后再 Mark / Commit |
| Rebalance | 旧成员还在处理,分区已经给了新成员 | 处理要能随时停;业务幂等;能用协作分配就用 |
| Commit 失败 | 处理成功,提交被拒或超时,重启后从旧位置重读 | 重试提交;用业务唯一键挡住第二遍 |
| 生产端重试 | Broker 已写,Producer 没收到 ack,再发一条新 offset | 幂等 Producer + 合理的 acks |
| 不同组订阅 | 两套系统各消费一遍同一份日志 | 这是设计如此。不要用同组硬凑「只处理一次」 |
组协议里的 session / heartbeat,和 Sarama 的 ConsumerGroupSession 不是同一个东西。前者是「我还活着,别把分区收回去」的租约;后者是客户端把一次分区归属包成的对象。心跳断了,Coordinator 认为你死了,分区易手,旧 session 的 Cleanup 会被调用。session.timeout 太短,GC 或一次稍慢的 process() 就会误杀;太长,真挂了要等很久才有人接盘。
Java 客户端还有 max.poll.interval.ms:两次 poll 间隔太久,组会认为你卡死,即使心跳还在。Sarama 的心跳走后台 goroutine,长处理不一定先死在心跳上,但仍会被 session.timeout / rebalance.timeout 盯着。所以不要在 ConsumeClaim 里做没有超时的重活;重活可以丢到自己的 worker,位移提交却必须和「这条已经产生了不可逆副作用」对齐。这正是多线程消费容易把 offset 交早的原因。
水位线:消费者究竟能读到哪里
leader 接受写入之后,消息不会立刻对 Consumer 可见。可见范围由几条水位共同决定。

- LEO(Log End Offset):这条副本下一条要写的位置。上图里是 15。
- HW(High Watermark):ISR 里所有副本都已经追上的位置。Consumer 只能读到 HW 之前,上图是 8。8 到 14 在 leader 本地已经落了,但还没被 ISR 认完,对消费者不存在。
- Log Start Offset:这条日志还留着的最早 offset。协议里有时叫 low watermark。retention(
log.retention.hours/log.retention.bytes)和 compaction 往前推它,更早的消息物理上没了。 - LSO(Last Stable Offset):开了事务之后才重要。
read_committed的消费者只能读到已提交事务,LSO 往往比 HW 更靠前。未提交事务挡在后面,避免读到会回滚的数据。
旧笔记里把 HW 写成「所有 follower 都同步完才往前挪」,太紧了。掉出 ISR 的副本不会拖住 HW。否则一台落后的盘就能把整个分区的读取卡住。
这也是 acks=1 会丢消息的几何原因:leader 的 LEO 已经到 15,HW 还在 8,此时 leader 挂了,新 leader 从 ISR 里选,HW 之后的那截在新 leader 上可能根本不存在。
副本、ISR,以及为什么数据面不用 Raft
复制因子(replication.factor)是创建 Topic 时定的,包含 leader 自己。rf=3 就是 1 个 leader + 2 个 follower。它不能大于 Broker 数;越大越抗丢,也越吃磁盘和网卡。
- Leader:该分区唯一的写入口,默认也是读入口
- Follower:主动 fetch leader 的日志,追上了才留在 ISR 里
- ISR:被认为「跟上了」的副本集合。判定窗口是
replica.lag.time.max.ms,不是固定条数差
Controller 只从 ISR 里选新 leader。unclean.leader.election.enable=false(默认就是 false)禁止 ISR 以外的副本抢权——否则可用性更好,但已经对消费者可见的数据可能被一笔勾掉。
面试里常把这套东西和 Raft 放在一起问。Kafka 的分区日志不用 Raft,原因很具体:
- 一个集群可以有几十万个 Partition。每个 Partition 拉一组 Raft,领导选举和心跳会把 Controller 和网络打爆
- 数据面要的是高吞吐追加,follower pull + ISR 比每条日志跑一遍 Raft 法定人数便宜
- 丢不丢、谁能当 leader,用 ISR 和
min.insync.replicas显式权衡就够了
元数据反过来很适合 Raft:全集群就一份,更新频率低,一致性要求高。所以 KRaft 用 Raft 管 Topic、分区、ISR 这些元数据,不管分区里的业务消息。老笔记里写的「ISR 和 Raft」应该这么拆,而不是把两者当成同一套协议。
Leader epoch(KIP-101 / KIP-320)补的是副本截断:新 leader 带着更大的 epoch 上任,follower 按 epoch 对齐,避免老 leader 未提交的尾巴把已经提交的数据盖掉。这是副本一致性真正精细的地方,不是再引入一套 Raft。
消息会在哪一层丢掉
三条路径。
Producer。 send 之后不管回调;acks=0;retries=0;消息太大或超时被拒,应用又没处理错误。对策是带回调发送、打开幂等、acks=all,瞬时错误重试,永久性错误(消息过大、鉴权失败)不要死循环。
Broker。 日志先落 page cache,刷盘是异步的。单机宕机本身就会丢还没刷下去的字节,所以可靠性从来不指望一台 Broker,而指望 ISR。acks=1 加 leader 崩溃、或者打开 unclean leader election,都会把已经「看起来成功」的消息弄没。
Consumer。 自动提交交早了,或者多线程里先提交再处理。消息在 Kafka 里还在,业务侧却再也不会读到——这是最常见的「丢」。关自动提交,副作用完成后再交 offset。
一组比较稳的底线:
acks=all
enable.idempotence=true
replication.factor=3
min.insync.replicas=2
unclean.leader.election.enable=false
enable.auto.commit=false
还要满足 replication.factor > min.insync.replicas,否则少一台 Broker,写入会直接拒绝,可用性没了。
三种交付语义
Kafka 对 Producer / Consumer 的承诺只有三种:
- At most once:可能丢,不会重复。Producer 不重试,或 Consumer 先提交再处理
- At least once:不丢,可能重复。这是默认,也是打开重试、先处理再提交之后的自然结果
- Exactly once:不丢也不重复
Exactly-once 在 Kafka 里不是一句开关。它拆成两截:
- 幂等 Producer:单分区、单 Producer 生命周期内的重试去重。Broker 多记一些字段,用空间换一次判断。
- 事务:把「读一批、写到另一个 Topic、提交这批 offset」收成一次原子提交。典型场景是 Kafka → Kafka 的处理链路,Kafka Streams 的 EOS 也走这条路。
read_committed的下游看不到未提交事务。
事务过不了这个边界:你的消费者读完 Kafka,写的是 MySQL、Redis、下游 HTTP。Kafka 无法和外部系统共一个事务。这时候「exactly-once」的落点是业务幂等——订单号唯一索引、去重表、把更新做成可重入。协议层能帮你的,只是少制造一些重复;少制造不等于零。
所以实践上几乎总是:按 at-least-once 来设计,用幂等把重复消化掉。只有整条链路都封闭在 Kafka 里时,才值得上事务。
为什么快
典型量级是单机每秒数十万条、延迟毫秒级。快不是因为「用了内存队列」,而是几件和磁盘、操作系统有关的选择叠在一起。
顺序写。 Partition 是纯追加,磁头不用来回寻道。机械盘的顺序写和 SSD 的差距,远小于随机写。
零拷贝。 传统读盘再发给网卡,数据要进内核、进用户态、再进 socket,来回拷。Kafka 在不需要转换格式时走 sendfile,page cache 直接到网卡。
flowchart LR
subgraph traditional [传统]
D1[磁盘] --> K1[内核]
K1 --> U[用户态]
U --> S1[Socket]
S1 --> N1[网卡]
end
subgraph zerocopy [零拷贝]
D2[磁盘] --> K2[内核 / page cache]
K2 --> N2[网卡]
end
Page cache,而不是 JVM 堆。 Broker 自己几乎不缓存消息体,把热数据交给操作系统。堆可以相对小,重启也不等于缓存冷启动——页还在 OS 里。消费者追上最新进度时,读的常常是还没刷盘、已经在 page cache 里的字节。
分区并行。 吞吐按 Partition 横向加,组内每个成员啃自己的那一份。
分段。 一个 Partition 不是单个无限大文件,而是一段段 segment,每段一对 .log + .index。

Topic-0
├── 00000000000000000000.log
├── 00000000000000000000.index
├── 00000000000000005000.log
└── 00000000000000005000.index
删过期数据是删文件,不是改一个巨大文件的头部;按 offset 查,先定位 segment,再走稀疏索引。
批量。 Producer 把多条收成一个 batch 再发出去。batch.size 和 linger.ms 是吞吐和延迟的同一根绳子:多等几毫秒,网卡和磁盘的单次开销摊到更多消息上。
压缩也是这根绳子上的一环。Producer 按 batch 压,Broker 默认原样存、原样转,Consumer 按 batch 头上的算法解。Broker 只在两件事发生时重压:集群配置了另一种算法,或要做消息格式转换。
| 吞吐 | 压缩比 | 备注 | |
|---|---|---|---|
| LZ4 | 最高 | 中 | 要吞吐时的默认选项 |
| Snappy | 高 | 偏低 | 老配置里最常见 |
| zstd | 中高 | 最高 | CPU 够、带宽紧时更合适 |
| gzip | 最低 | 高 | CPU 贵,现在很少作为首选 |
Producer 和 Consumer 的 CPU 都富余、带宽却紧,开压缩才划算。两边都已经 100% 了,再压是给延迟加戏。
生产里这组旋钮够用:
# Producer
batch.size=16384
linger.ms=5
compression.type=lz4 # 或 zstd
# Broker
num.network.threads=3
num.io.threads=8
# Consumer
fetch.min.bytes=1024
max.partition.fetch.bytes=1048576
具体数字要按消息大小和延迟预算改,没有一份能抄到所有 Topic 上的配置。
实践上我会怎么配
- 同一类处理进程共用一个
group.id;不同下游系统必须拆组。不要为了「别重复」把风控和报表塞进同一个组。 enable.auto.commit=false。处理出了不可逆副作用之后再 Mark /commitSync。- 生产端打开幂等;可靠性优先就
acks=all+min.insync.replicas=2+rf=3。 - 业务层必须幂等。消息 ID、订单 ID、事件 ID,落到 DB 唯一约束或 Redis SETNX。这是重复消费真正的最后一道门。
- 能用协作式分配就用,减少 rebalance 时的停顿。但不要幻想换个 assignor 就能 exactly-once。
- 整条链路都在 Kafka 里、真的需要原子「读-转化-写」,再上事务或 Kafka Streams EOS。
- 分区数按峰值并行度提前规划。Partition 对业务方来说只能加不能减;加了之后旧数据不会重切,hash 空间却变了,同一 key 的历史和新消息可能从此不在一个分区。
- 心跳和
poll间隔按最慢那次处理来留余量。一次process()里不要做没有超时的远程调用。
一句话收束:Producer 负责把消息可靠地追加到某个 Partition;Consumer 靠 poll 和 offset 拉自己被分到的那一段;组内用分区独占避免并发重复。崩溃和 rebalance 仍然会重放,重复消费控制的终局是业务幂等,不是某个客户端参数。