KEEL · 龙骨 · A CURRICULUM FOR THE AI ERA

00 · 提交日志不是队列 — keel 龙骨

这一章纠正一个前提。你后面所有的判断——acks 该设几、分区该给几个、积压该不该慌——都建立在「Kafka 里一条消息被消费之后会发生什么」这个前提上。前提错了,后面全错。

这一章纠正一个前提。你后面所有的判断——acks 该设几、分区该给几个、积压该不该慌——都建立在「Kafka 里一条消息被消费之后会发生什么」这个前提上。前提错了,后面全错。

现场

一个电商系统把订单事件写进 Kafka。下游有两个服务:推荐服务要把订单流喂给特征工程,风控服务要用同一批订单算实时风险。两个团队各自写了消费者,各自起了一个消费组,各自说「我这边收到的订单是全的」。

如果 Kafka 是队列,这个场景根本不可能成立——同一条消息被推荐服务取走,风控服务就再也看不到了,两个服务只能抢消息。但它成立。要理解为什么成立,得先看清 Kafka 到底把「读到哪了」记在了什么地方。

两套系统对照

flowchart TD
  subgraph Q["队列语义(消息被取走即删除)"]
    direction TB
    QP["生产者"] --> QQ["队列 m1 m2 m3 m4 m5"]
    QQ --> QC1["消费者A 取走 m1 m2 m3"]
    QQ -. "已被取走的 m1..m3 消费者B 看不到" .-> QC2["消费者B 只能拿 m4 m5"]
    QC1 --> QDEL["broker 删除 m1 m2 m3"]
  end

  subgraph L["Kafka 日志语义(追加 + 每个组各自记位点)"]
    direction TB
    LP["生产者<br/>orders-02"] --> LOG["orders-02-0 追加日志<br/>offset 0..1675 全部留在盘上"]
    LOG --> GA["消费组 ga<br/>位点 1676"]
    LOG --> GB["消费组 gb<br/>位点 1676"]
    LOG -. "保留多久由 retention 决定,不由消费决定" .-> RET["retention.ms"]
  end

  style QP fill:#e8f0fe,stroke:#3b5bdb,color:#1a1a2e
  style QQ fill:#e8f0fe,stroke:#3b5bdb,color:#1a1a2e
  style QC1 fill:#e8f0fe,stroke:#3b5bdb,color:#1a1a2e
  style QC2 fill:#e8f0fe,stroke:#3b5bdb,color:#1a1a2e
  style QDEL fill:#ffe3e3,stroke:#c92a2a,color:#1a1a2e
  style LP fill:#e6fcf5,stroke:#0ca678,color:#1a1a2e
  style LOG fill:#e6fcf5,stroke:#0ca678,color:#1a1a2e
  style GA fill:#fff4e6,stroke:#e8590c,color:#1a1a2e
  style GB fill:#fff4e6,stroke:#e8590c,color:#1a1a2e
  style RET fill:#f3f0ff,stroke:#7048e8,color:#1a1a2e

两条链路的差别不在吞吐,而在**「读到哪了」这件事记在谁身上**:队列把它记在 broker 的消息删除上,Kafka 把它记在每个消费者自己的位点里。图里 ga 和 gb 是两个互不相干的指针,指向同一份日志的不同位置——本机上它们恰好都指向末尾。

精确定义

offset 不是消息 ID,是消息在单个分区里的序号,从 0 开始单调递增。分区内的 offset 连续、不重复、不跳号(压缩过的 topic 例外,见下)。一条消息在全局没有唯一 offset——orders-02-0 的 offset 3 和 orders-02-1 的 offset 3 是两条不同的消息。

消费位点(committed offset) 是某个消费组在某个分区上「下一条要读的位置」。它不是消息上的字段,而是消费组写回 Kafka 的一条记录,存在内部 topic __consumer_offsets 里。

滞后(lag) = 分区末端 offset(LOG-END-OFFSET)− 消费组位点(CURRENT-OFFSET)。它衡量的是这个组还有多少条没读完,不是集群的健康度。

一条消息在被消费后不会被删除。它留在段文件里,直到 retention 规则(retention.ms / retention.bytes)或日志压缩把它清掉。第 08 章的段滚动实验会看到这个删除动作实际发生。

一次完整运行

本机实验(原始输出见 lab/evidence/kafka-engineering/00-log-not-queue.txt)。orders-02 是一个 3 分区、3 副本的 topic,末端 offset:

orders-02:0:1676
orders-02:1:1640
orders-02:2:1617

三个分区加起来 1676 + 1640 + 1617 = 4933 条。现在用两个全新的消费组 ga 和 gb,各自从最早(--from-beginning)读一遍这个 topic:

group ga -> Processed a total of 4933 messages
group gb -> Processed a total of 4933 messages

两个组各自读到了完整的 4933 条。这不是巧合,是设计的直接结果:读动作不改写日志,只是各自把一个数字往前推。再看两个组各自的位点:

--- group ga ---
ga   orders-02   0   1676   1676   0
ga   orders-02   1   1640   1640   0
ga   orders-02   2   1617   1617   0
--- group gb ---
gb   orders-02   0   1676   1676   0
gb   orders-02   1   1640   1640   0
gb   orders-02   2   1617   1617   0

四列分别是 CURRENT-OFFSET(位点)、LOG-END-OFFSET(末端)、LAG(滞后)、以及是否有活跃成员。两个组的位点都到了末端,lag 都是 0,但它们是两条独立的记录——一个组读完了,另一个组完全不受影响。这一栏里 CURRENT-OFFSET 和 LOG-END-OFFSET 相等就代表「这个组追平了」,注意这里的 3 个分区、6 行位点,全部来自这一份日志。

位点到底存在哪

__consumer_offsets 这个 topic 就是答案:

Topic: __consumer_offsets  PartitionCount: 50  ReplicationFactor: 3
Configs: compression.type=producer,min.insync.replicas=2,cleanup.policy=compact,...

几个字段值得逐个读:

理解这一层,第 02 章讲 acks 时你就能看明白:生产者的 acks 管的是业务消息,位点的持久化是另一套(消费组提交时走 __consumer_offsets 的副本机制)。两条路径是分开的。

这套设计换来了什么

回放。位点是一个数字,想重读就把它往前调。生产事故里最常见的操作是「这个组的逻辑改了,把位点倒回 24 小时前重跑一遍」,--reset-offsets --to-earliest 做的就是这个(第 04 章会实测)。队列做不到这一点,消息删了就是删了。

多个消费组互不干扰。推荐和风控各读各的,谁的消费慢了都不会拖累对方。这也是「同一个 topic 订阅两个下游」的标准做法。

读放大不额外占空间。N 个消费组读同一份日志,磁盘上还是一份数据。

放弃了什么

按消息级别的确认与删除。RocketMQ、RabbitMQ 的「这条我处理完了,删掉」在 Kafka 里没有对应物。你能做的只有推进位点。想让「个别失败消息」留在那里不前进,得把整个分区的位点卡住——这就是为什么死信队列在 Kafka 里要用另一个 topic 自己搭。

灵活路由。Kafka 投递到哪个分区由 key 哈希或轮询决定,broker 不做基于内容的路由。想按业务字段分发,得在生产者侧自己做,或者用 Streams 再写一个 topic。

单条消息的 TTL。retention 是 topic 级别(或分区级别)的时间/大小,不能给某一条消息单独设过期。8 秒的定时消息、24 小时过期的验证码,这些在 Kafka 里都得自己算,或者换 RocketMQ(见 《可扩展性》05 章)。

分区数决定并行上限。一个消费组最多有「分区数」个消费者在干活,第 4 个消费者起来会闲置。这是队列模型没有的约束——队列可以无限加消费者抢同一条队列。

顺带说一个反直觉的点

消费组「读完」并不代表消息没了。ga 读完之后,orders-02 的数据还在盘上,末端 offset 也还是 1676。它什么时候消失,取决于 retention.ms(本机 orders-12 上设的是 15 秒,第 08 章会看到段文件真的被删)。这两个概念——「我读到哪」和「数据还能留多久」——必须分开记,混在一起就会出现「消费完了,为什么还能重读」和「明明读过了,怎么数据没了」两种相反的困惑。


生产边界

动手

  1. 给 orders-02 再建一个消费组 gc,从最早读一遍,确认它也读到 4933 条,并且 ga/gb 的位点没有任何变化。
  2. 用 ./kafka-cli.sh offsets --bootstrap-server localhost:9092 --topic orders-02 记下末端 offset,生产 10 条新消息,再跑一次,确认只有变化的分区 offset 增长。
  3. 断言题:把 gb 的位点用 --reset-offsets --to-earliest 倒回 0,然后 --describe 看它的 LAG 是否等于该分区的 LOG-END-OFFSET。

自测

  1. 同一条消息被 ga 读到之后,gb 为什么还能读到?读到哪这件事记在哪个组件上?
  2. offsets --topic 打印的三个数字(分区、末端 offset)里,哪个是「这个分区一共有多少条消息」?为什么不能用它减 1 当全局消息数?
  3. __consumer_offsets 用的是 cleanup.policy=compact 而不是按时间删除,这么做解决了什么问题?如果改成按时间删除会出什么事故?
  4. 一个消费组 lag 长期为 0,能不能说明集群是健康的?举一个 lag=0 但数据已经出问题的场景。
  5. 「消费完消息就没了」这个直觉在什么情况下会真的成立?(提示:和第 08 章的 retention 联系起来想。)

↓ 下一步:01 章 · 存储模型

进入 keel 阅读