返回归档
🕹️工程与设计

Kafka 生产者、消费者与重复消费

Kafka 是按 Partition 切开的提交日志。组内靠分区独占避免并发抢同一条消息;崩溃、rebalance、提交失败仍会让同一条消息被处理两遍。协议层和业务层要分开看。

文章目录

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。后文会回到这一点。

关系和基数可以先记在这张图里:

Kafka 核心概念:Producer、Topic、Partition、Replica 与 Consumer Group

几个容易混的数量关系:

  • 一个 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 同时喂给风控、报表、通知,靠的就是这个,不是什么额外的广播模式。

生产者怎么写进去

一次写入大致是四步。

  1. 应用构造一条 ProducerRecord:topic、key、value,再加上 timestamp、headers。
  2. 分区器选定目标 Partition。有 key 时通常按 hash(Java 客户端是 murmur2)路由,同一个 key 会进同一个 Partition,分区内有序。无 key 时,2.4 之后默认是 sticky partitioner:先粘在一个分区上把 batch 攒满,再换下一个,吞吐比纯轮询好。
  3. 消息发给该 Partition 的 leader。
  4. 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=3min.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
}

注意两件实现细节:

  1. ConsumeClaim 必须把 claim.Messages() 读干,或者立刻退出让 session 结束。卡在中间不读也不返回,Sarama 的心跳仍可能在后台续着,但分区收不回去、rebalance 也会卡住。
  2. 处理失败就不要 Mark,也不要 Commit。Kafka 没有 NACK,不提交就是「下次还从这里开始」。offset 是分区上的一个数字,后面的消息一旦被提交,失败的那条就被跳过了。

组内为什么看起来不重复

多个进程共用同一个 group.id,就是同一个 Consumer Group。对同一个 Topic 的同一个 Partition,组内同一时刻只有一个成员拥有消费权。这是 Kafka 避免组内并发重复的全部手段。

Consumer Group:Coordinator 按分区分配消费权,每个成员维护自己的 offset

六个 Partition、三个 Consumer,理想情况是每人两个。再加第四个、第五个、第六个还能继续摊薄;第七个开始就会有人空转——并行度的上限是 Partition 数,不是 Consumer 数

图里写在 Consumer 身上的「Offset: 42」是简化。真实进度是每个 (group, topic, partition) 一份。一个人吃两个分区,就有两条 offset。

再强调一次反面:

  • 同一个 Partition 不能被同组两个 Consumer 同时拉
  • 一个 Consumer 可以同时拉多个 Partition
  • 不同组会各自完整地消费同一份 Topic,这是特性,不是 bug
  • 组与组之间的 offset 完全独立

顺序同样只在分区里成立。要让同一个订单的事件保序,就把 order_id 当 key,让它们进同一个 Partition。跨分区没有全局顺序。

Rebalance:重复处理的高发区

组成员变了,分区就要重分。常见触发:

  1. 有人加入
  2. 有人主动离开
  3. 有人崩溃或心跳超时
  4. Topic 的分区数变了
  5. 订阅关系变了

旧协议是 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=0retries=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 里不是一句开关。它拆成两截:

  1. 幂等 Producer:单分区、单 Producer 生命周期内的重试去重。Broker 多记一些字段,用空间换一次判断。
  2. 事务:把「读一批、写到另一个 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 → Partition → Segment → log/index

Topic-0
├── 00000000000000000000.log
├── 00000000000000000000.index
├── 00000000000000005000.log
└── 00000000000000005000.index

删过期数据是删文件,不是改一个巨大文件的头部;按 offset 查,先定位 segment,再走稀疏索引。

批量。 Producer 把多条收成一个 batch 再发出去。batch.sizelinger.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 上的配置。

实践上我会怎么配

  1. 同一类处理进程共用一个 group.id;不同下游系统必须拆组。不要为了「别重复」把风控和报表塞进同一个组。
  2. enable.auto.commit=false。处理出了不可逆副作用之后再 Mark / commitSync
  3. 生产端打开幂等;可靠性优先就 acks=all + min.insync.replicas=2 + rf=3
  4. 业务层必须幂等。消息 ID、订单 ID、事件 ID,落到 DB 唯一约束或 Redis SETNX。这是重复消费真正的最后一道门。
  5. 能用协作式分配就用,减少 rebalance 时的停顿。但不要幻想换个 assignor 就能 exactly-once。
  6. 整条链路都在 Kafka 里、真的需要原子「读-转化-写」,再上事务或 Kafka Streams EOS。
  7. 分区数按峰值并行度提前规划。Partition 对业务方来说只能加不能减;加了之后旧数据不会重切,hash 空间却变了,同一 key 的历史和新消息可能从此不在一个分区。
  8. 心跳和 poll 间隔按最慢那次处理来留余量。一次 process() 里不要做没有超时的远程调用。

一句话收束:Producer 负责把消息可靠地追加到某个 Partition;Consumer 靠 poll 和 offset 拉自己被分到的那一段;组内用分区独占避免并发重复。崩溃和 rebalance 仍然会重放,重复消费控制的终局是业务幂等,不是某个客户端参数。