Skip to content
Charles Shao
Go back

Kafka 核心原理精讲:从一条日志到分布式流平台,把关键机制与取舍讲透

views

在广告系统里,一次曝光(impression)、一次点击(click)、一次转化(conversion)都是一个事件。我关心的问题是:一天几十上百亿条这样的事件,如何可靠地、按序地、能被多方同时消费地从投放端流向计费、报表、风控、特征管道与实时数仓。Apache Kafka 最初在 LinkedIn 诞生,正是为了统一处理这类高吞吐、需持久化、要被多方订阅与重放的事件流——广告事件流只是它最典型的应用场景之一。

本文以 Jay Kreps 提出的「日志为中心」抽象为线索,把 Kafka 从一条 append-only 日志推演到分布式流平台,并沿途标出我认为真正决定生产行为的机制与取舍。

TL;DR

Table of contents

Open Table of contents

1. 内核:Kafka 是一个”分布式提交日志”

抛开所有花哨的名词,我把 Kafka 的心脏归结为一个极简的数据结构:一条只追加(append-only)、按写入顺序排列、每条记录带一个单调递增编号(offset)的日志。

Kafka 分区 append-only 日志示意图。中间是一个 Partition,画成从左到右一串带编号的格子:offset 0、1、2、3、4、5,最右边是 offset 6,标注为 LEO(Log End Offset,下一条消息将写入的位置)。左侧的 Producer 用箭头指向日志末尾(LEO),标注『只能 append 到末尾』,说明写入只发生在日志尾部。右侧有两个消费者:Consumer A 的指针停在 offset 2(committed offset = 2),Consumer B 的指针停在 offset 5(committed offset = 5),说明不同消费者在同一条日志上各自维护独立的读取位置。底部旁注:读不删除——消费只是推进各自的 offset,不同消费者各读各的,并且可以把 offset 重置回任意历史位置进行重放。

就这么一个”账本”,配合三个关键决定,就长成了 Kafka:

  1. 写入只追加到末尾 —— 顺序写磁盘,快到接近写内存。
  2. 读不删除、按策略留存 —— 消费不是”取走”,而是”读到某个位置”。同一份日志因此能被 N 个互不相干的消费者各读各的。
  3. 把这份日志切片、复制、分散到多台机器 —— 切片得到 Partition(拿到并行与扩展),复制得到 Replica(拿到高可用)。

一句话:Kafka = 把一个”可持久化、可重放的日志”,做成了分布式、可水平扩展、多方可订阅的系统。后面所有概念,都是围着这句话展开的。

这也解答了一个经典追问——“Kafka 到底是 MQ、数据库还是流处理引擎?“我的结论是:它本质是日志,MQ / 存储 / 流处理只是这份日志的不同用法。这个”日志为中心(log-centric)“的世界观,是 Kafka 与传统”队列为中心”的 MQ 最根本的分野(§11 会展开)。

1.1 一条 Record 长什么样

日志里的每个元素叫 Record(记录/消息),它不是裸字符串,而是一个结构:

字段含义
Key(可空)决定分区路由与压缩语义的键(如 user_idorder_id
Value真正的负载(payload),如一条曝光事件的 JSON / Avro / Protobuf
Timestamp事件时间或写入时间
Headers可选的元数据键值对(如 trace id、schema 版本)
Offset由 broker 在写入分区时分配的、分区内唯一且递增的编号

Key 是被低估的关键字段:它同时决定”这条消息进哪个分区(→ 顺序性)“和”日志压缩时以什么为主键去重”。我在生产里见过的最常见性能坑之一,就是 Key 分布不均导致的热点分区(hot partition)(§12)。

2. 从”日志”到”分布式事件流平台”

官方把 Kafka 定义为 distributed event streaming platform(分布式事件流平台),并强调三大核心能力:发布/订阅、持久化存储、流式处理;再算上打通外部系统的 Connect,日常打交道的其实是下面这四类(前两类是本文重点):

能力含义典型载体
① 发布/订阅生产者发、消费者订阅一条条消息流Producer / Consumer API(当 MQ 用)
② 持久化存储消息按容错方式落盘留存,可重放分区日志 + 副本(当”可回放事件存储”用)
③ 流式处理在消息流动时做转换 / 聚合 / 连接Kafka Streams / ksqlDB / Flink
④ 系统集成与数据库、对象存储等打通的连接器Kafka Connect(Source / Sink)

在 AdTech 里这几类能力常常同时用

3. 核心概念全景

先用一张总览图建立全局空间感——上半部分是各核心概念之间的关系与数量基数(cardinality),下半部分是生产/消费的部署视角——再逐个精确拆解:

Kafka 架构总览图,分 A、B 两部分。A 部分是核心概念的实体关系图,用箭头上的数字标注数量基数(1 = 恰好一个,0..* = 零到多个,1..* = 一到多个):一个 Producer 写入 1 个 Topic,一个 Topic 有 0 到多个 Producer;一个 Topic 有 0 到多个 Consumer;一个 Consumer 恰好属于 1 个 Consumer Group;一个 Topic 切分为 1 到多个 Partition;同一个 Consumer Group 内一个 Partition 只对应 1 个 Consumer,而一个 Consumer 可从 0 到多个 Partition 拉取消息;一个 Cluster 含 1 到多个 Broker,一个 Broker 属于 1 个 Cluster;一个 Partition 有 1 到多个 Replica,每个 Replica 落在某 1 个 Broker 上,其中 1 个是 Leader Replica、其余 0 到多个是 Follower Replica。B 部分是部署视角:左侧多个 Producer 以红色虚线 push 写入中间 Kafka Cluster 里的多台 Kafka Broker(每台存有 P0 / P1 / P2 等分区);右侧多个 Consumer(其中两个组成一个 Consumer Group)以绿色虚线 pull 拉取消息;集群的元数据与 Leader 选举由内置的 KRaft Controller Quorum 仲裁。

上图 B 部分的元数据管理画的是 KRaft(Kafka 3.3+ 生产可用、4.0 起彻底移除 ZooKeeper,见 §10);更早的版本这里是一个外部的 ZooKeeper 集群负责注册与选主,但除此之外的所有关系(Topic / Partition / Replica / Consumer Group 等)完全不变。

3.1 Topic(主题)

消息的逻辑分类 / 命名频道。生产者往 topic 发、消费者从 topic 订阅,类比”一个具名的事件流”(如 ad-impressionsad-clicks)。Topic 是逻辑概念,真正存数据的是它下面的 Partition。

3.2 Partition(分区)—— 分布式的最小单位

一个 topic 被切成一个或多个 partition,每个 partition 就是第 1 节那条 append-only 有序日志。 它是 Kafka 分布式、并行、扩展、顺序、容量的共同基本单位,而这些属性彼此牵制:

这条”只保证分区内有序”是 Kafka 最重要、也最常被误解的性质。要顺序,就得让相关消息落到同一个 partition(靠 key);要并行,就得把消息摊到多个 partition——二者天然冲突,这正是顺序性与吞吐取舍的根。分区数怎么定,见 §7。

3.3 Offset(偏移量)

每条消息在其 partition 内唯一、单调递增的编号。 三个要点:

正因进度归消费者自己管,Kafka 才天然支持重放(把 offset 重置到过去)与多方独立消费(各记各的 offset)。

3.4 Broker 与 Controller

一台 Kafka 服务器就是一个 broker,负责存 partition、处理读写请求。多个 broker 组成集群,partition 及其副本分散其上。集群中有一个特殊角色 Controller(控制器),负责分区 Leader 选举、副本状态管理等元数据工作(KRaft 模式下由 controller quorum 通过 Raft 承担,见 §10)。

3.5 Producer(生产者)

把消息写进 topic 的客户端。两个核心决策:这条消息进哪个 partition(分区策略)与多可靠地确认写入acks)。这两个旋钮几乎决定了顺序、吞吐、可靠性的全部取舍,是第三篇的主角。

3.6 Consumer / Consumer Group(消费者 / 消费者组)

于是:想让多个系统各拿一份完整数据(广播)→ 用不同 group;想让一个系统内并行分担负载 → 同一 group 内加消费者(上限 = partition 数)。为什么 Kafka 选 pull 而非 push,见 §5.1。

一图看清”组内消费者数与分区数”的关系:

Kafka 消费者组与分区对应关系图。正中是一个 Topic,内含 Partition 0 / 1 / 2 三个分区;四周环绕四个消费者组,演示"组内消费者数量与分区数量"如何决定消费方式。Consumer Group 1 只有一个消费者 C11,用红色箭头独自消费全部三个分区。Consumer Group 2 有两个消费者,C21 分到一个分区、C22 分到两个分区。Consumer Group 3 有三个消费者 C31 / C32 / C33,恰好一人一分区,达到最高并行度。Consumer Group 4 有四个消费者,但分区只有三个,C41 / C42 / C43 各分到一个分区,多出来的 C44 因为没有分区可分而处于空闲(idle)状态。四个组各自独立维护自己的 offset、互不影响,所以同一条消息可被多个组分别消费;而在同一个组内,一个分区只会被一个消费者消费,因此一个组的消费并行度上限就等于分区数。

3.7 Replica / ISR / High Watermark —— 可见性与一致性的边界

这套机制用一张时序图最清楚:

Kafka High Watermark 与副本复制时序图。参与方从左到右:Producer、Leader(Broker 1)、两个 ISR 内的 Follower(Broker 2、Broker 3)、Consumer。流程:① Producer 以 acks=all 写入消息;② Leader 追加到本地日志,LEO 加 1,此时该消息在 Leader 尾部但尚未提交;图中旁注:此刻消息在 Leader 上但 High Watermark 未推进,对 Consumer 不可见;③④ 两个 Follower 分别向 Leader fetch 拉取这条新消息完成复制;⑤ Leader 发现所有 ISR 都已复制,于是推进 High Watermark;⑥ Leader 才向 Producer 返回 ack,表示 HW 已覆盖该消息、消息不会丢;⑦ Consumer 向 Leader fetch 消费;⑧ Leader 只返回 offset 小于等于 HW 的消息,HW 与 LEO 之间尚未被全部 ISR 复制的尾部消息对 Consumer 不可见。

为什么消费者只能读到 HW? 因为 HW 之前的数据已被所有 ISR 复制,即使此刻 leader 宕机、从 ISR 选出新 leader,这些数据也一定还在——消费者看到的永远是”已确保不会因换主而消失”的部分。这就是 Kafka 在”可用性”与”一致性”之间划下的那条线。

那 HW 与 LEO 之间的数据呢? Leader 宕机时,这段”已写在旧 Leader 上、但还没被全部 ISR 复制”的尾部可能丢失:新 Leader 只从 ISR 选出,其日志通常只到旧 HW;旧 Leader 上独有的尾巴不再被承认。因此:≤ HW 安全(换主后仍在、可消费);HW~LEO 可能丢(换主后不再承认)。Producer 若尚未收到 acks=all 的 ack,换主后应重试;消费者则从未见过这段数据,也就不会出现”先读到、再消失”。

4. 消息模型:为什么 Kafka 选了发布-订阅

传统 MQ 有两种模型。我关心的不是名词本身,而是 Kafka 如何用一套机制把两者统一:

队列模型与发布-订阅模型对比图。上半部分是队列模型:一个 Producer 把消息发进一个 Queue,Queue 把消息分摊给两个消费者——Consumer 1 拿到 m1、m3……,Consumer 2 拿到 m2、m4……,即一条消息只被其中一个消费者取走。下半部分是 Kafka 采用的发布-订阅模型:一个 Producer 把消息发进一个 Topic,Topic 把全量消息 m1..mN 分别广播给 Consumer Group A 和 Consumer Group B,两个组各自都能收到完整的一份数据。底部旁注:Kafka = topic(广播给多个 Group)+ Group 内负载均衡(分摊);当只有一个 Group 时就退化成队列模型,因此 pub-sub 在功能上向下兼容队列模型。

队列模型:消息进队列,一条只被一个消费者取走。适合任务分发,但想让”多个系统各拿一份完整数据”就很别扭(得为每个消费者复制一份)。

发布-订阅模型(Kafka):以 topic 为载体,发布一条,所有订阅方各自都能收到。当一个 topic 只有一个订阅组时,它退化成队列模型——所以 pub-sub 在功能上向下兼容队列模型。

Kafka 用 topic(广播)+ consumer group(组内负载均衡) 一举同时表达两种语义:不同 group = 广播(各拿全量);同 group 多消费者 = 队列式负载均衡(分摊)。

RocketMQ 的消息模型与 Kafka 基本一致;差异之一是 Kafka 没有”队列”这个概念,与之对应的是 Partition

5. 消费者组与再平衡(Rebalance)

消费者可随时加入或退出 group。当组成员变化(或订阅的分区数变化)时,分区所有权会被重新分配给存活消费者——这个过程叫 rebalance(再平衡)。

消费者组再平衡(rebalance)两种策略对比图,分上下两条泳道。上条是 Eager rebalance(默认 RangeAssignor):从『消费中(C1 持 P0、P1,C2 持 P2、P3)』出发,成员变化时全体消费者先 revoke 全部分区,进入一个红色高亮的『全组停摆 stop-the-world 窗口』,再重新入组统一分配,最后才恢复消费(含新成员 C3)——中间整组暂停。下条是 Cooperative rebalance(CooperativeStickyAssignor):同样从『消费中』出发,但只撤回需要迁移的少数分区、其余分区照常消费,随后增量分配被撤回的那几个分区并直接恢复,全程没有全组停摆。底部旁注:成员增减或订阅变化会触发 rebalance(组内一分区仍只归一个 Consumer);Eager 有全组停摆窗口,Cooperative 只挪一小部分、其余不停;Kafka Streams 默认 cooperative,普通消费者需显式配置。

它带来高可用与弹性伸缩,但也有代价,分两种策略:

Rebalance 期间消费暂停会造成处理延迟。对我做过的实时反作弊一类链路,这是要重点规避的抖动源。常见做法:控制单批处理耗时、合理设置 max.poll.interval.ms 与心跳,避免消费者被误判”假死”而触发重平衡。

5.1 为什么 Kafka 是 pull(拉)模型

很多 MQ 是 broker 主动 push 给消费者,Kafka 却让消费者主动 pull。这不是随意选择,而是一组深思熟虑的取舍:

维度Pull(Kafka)Push(传统 MQ)
流控(backpressure)消费者按自己节奏拉,天然不会被打爆broker 需感知消费者速率,慢消费者易被压垮
批量消费者一次拉一批,摊薄网络/处理开销逐条推难以高效批量
重放消费者可任意重置 offset 回放历史消息推走即难回放
代价无数据时会空轮询(用 long-poll fetch.max.wait.ms 缓解)有数据即时下发,低延迟

pull 模型是”dumb broker + smart client”哲学(§11)的直接体现:broker 不追踪、不推送、不管消费进度,把控制权全交给客户端。代价是消费者要自己处理”没数据时怎么等”,收益是 broker 极简、可扩展、可回放。

6. 上手:把概念跑起来(KRaft + CLI + Java)

原理讲了这么多,最快的内化方式仍是亲手发一条、消一条。下面这套 demo 我按 KRaft 单节点写好,可直接复制运行,全程无 ZooKeeper。

6.1 起一个单节点集群(KRaft,无 ZooKeeper)

官方镜像自 3.7 起默认以 KRaft 模式运行(controller + broker 合一),一条命令即可:

docker run -d --name kafka -p 9092:9092 apache/kafka:3.9.0
# 容器内的 CLI 脚本在 /opt/kafka/bin/ 下

6.2 建 topic + 控制台收发(先跑通链路)

# 建一个 3 分区的 topic(单副本,仅本地体验用)
docker exec kafka /opt/kafka/bin/kafka-topics.sh --create \
  --topic ad-events --partitions 3 --replication-factor 1 \
  --bootstrap-server localhost:9092

# 生产:以 : 分隔 key 与 value,key 决定进哪个分区
docker exec -it kafka /opt/kafka/bin/kafka-console-producer.sh \
  --topic ad-events --bootstrap-server localhost:9092 \
  --property parse.key=true --property key.separator=:
# 然后逐行输入:  user-42:{"type":"impression"}

# 另开一个终端消费(--from-beginning 从头读,直观感受"读不删除")
docker exec -it kafka /opt/kafka/bin/kafka-console-consumer.sh \
  --topic ad-events --from-beginning --bootstrap-server localhost:9092 \
  --property print.key=true

--from-beginning 的消费者关掉再开一个--group 的,会发现它接着上次的位点读;换个新 --group 又能从头再读一遍——这就是 §3.3 的 offset 与”多方独立消费 / 重放”在命令行里的直接体现。

6.3 Java Producer:acks=all + 幂等 + key 路由

Maven 只需 org.apache.kafka:kafka-clients(下面用 import static ProducerConfig.* 让配置项更短)。这段把三个关键旋钮一次拧到位:

Properties props = new Properties();
props.put(BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
props.put(KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
props.put(VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
props.put(ACKS_CONFIG, "all");              // 等 ISR 全部复制才算成功 → 不丢
props.put(ENABLE_IDEMPOTENCE_CONFIG, true); // 幂等:重试不产生重复、不乱序

try (KafkaProducer<String, String> producer = new KafkaProducer<>(props)) {
    // key = user_id:同一用户的事件落同一分区 → 该用户维度分区内有序
    producer.send(
        new ProducerRecord<>("ad-events", "user-42", "{\"type\":\"impression\"}"),
        (meta, ex) -> {
            if (ex != null) ex.printStackTrace();
            else System.out.printf("→ partition=%d offset=%d%n", meta.partition(), meta.offset());
        });
} // close() 会 flush 未发送完的批次

6.4 Java Consumer:手动提交 = at-least-once + 幂等消费

生产里最常用的组合:先处理、后提交(at-least-once),再用业务主键去重兜住重复(§9)。配置项同样用 import static ConsumerConfig.* 简写:

Properties props = new Properties();
props.put(BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
props.put(GROUP_ID_CONFIG, "billing");
props.put(KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
props.put(VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
props.put(ENABLE_AUTO_COMMIT_CONFIG, false);   // 关自动提交,改为处理完手动提交
props.put(AUTO_OFFSET_RESET_CONFIG, "earliest");

try (KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props)) {
    consumer.subscribe(List.of("ad-events"));
    while (true) {
        ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(500));
        for (ConsumerRecord<String, String> r : records) {
            String bizKey = r.key() + "-" + r.partition() + "-" + r.offset();
            if (seen.putIfAbsent(bizKey, Boolean.TRUE) == null) { // 幂等去重
                process(r);   // 真正的副作用:写库 / 计费
            }
        }
        consumer.commitSync(); // 处理成功后才提交 → 崩溃只会重复投递,绝不丢
    }
}

顺序很关键process()commitSync() 之前 → at-least-once(崩溃会重复、不会丢);若把两者调换(先提交再处理)就变成 at-most-once(崩溃会丢)。这正是 §9 那张交付语义表在代码里的样子。

别把这里的 seen 当真幂等:它只是个进程内的 Map,用来演示”去重发生在处理之前”——进程重启即失效,多实例之间也无法共享。生产里的幂等必须落到外部持久存储:数据库唯一键、INSERT ... ON CONFLICT、Redis SETNX、或专门的去重表(§9),这样才能真正兜住 at-least-once 带来的重复投递。

6.5 重放与观测:Kafka 的两个杀手级操作

# 重放:把 billing 组的位点重置到最早,重新消费全部历史(对账 / 重算)
docker exec kafka /opt/kafka/bin/kafka-consumer-groups.sh \
  --bootstrap-server localhost:9092 --group billing \
  --topic ad-events --reset-offsets --to-earliest --execute

# 观测:查看每个分区的消费滞后(LAG)——生产环境头号健康指标
docker exec kafka /opt/kafka/bin/kafka-consumer-groups.sh \
  --bootstrap-server localhost:9092 --describe --group billing

LAG = 分区可消费末端(工具里的 LOG-END-OFFSET 列,正常时约等于高水位 HW)− 该组已提交位点(CURRENT-OFFSET,即”还差多少没消费”。它是判断消费者是否跟得上生产速率、要不要加分区/消费者的第一诊断信号

7. 分区数:一个绕不开的容量决策

分区数(partition count)是建 topic 时最重要、又最容易拍脑袋的参数。它同时决定四件事,且每个方向都有代价:

分区太少分区太多
消费并行度受限(group 内消费者数 ≤ 分区数)元数据膨胀,Controller 选举 / 故障恢复变慢
单分区吞吐见顶,难以水平扩展每分区多套文件句柄与内存,broker 资源吃紧
热点 key 更易把单分区打满端到端延迟上升(更多分区要 flush / 复制)
——Producer 端每分区一个批次,攒批效率下降

Kafka 分区数取舍曲线图,横轴是分区数(由少到多),纵轴是相对量。绿色曲线表示『吞吐 / 并行度收益』,随分区数增加先快速上升、随后趋于饱和(收益递减);红色曲线表示『选举 / 恢复 / 延迟 / 文件句柄成本』,起初平缓、之后加速上扬,并在右侧某一点反超收益曲线。两线之间、绿色收益已接近饱和而红色成本仍较低的那一段被框为绿色的『sweet spot 合适区间』。曲线左端标注『太少:并行受限』,右端标注『太多:成本陡增』。底部黄色提示条给出反推公式:目标分区数 ≈ max(目标吞吐 / 单分区吞吐, 期望并行度)+ 适度余量,并强调分区数只增不减,宁可略多、别无脑堆到成千上万。整张图传达的核心是:收益会饱和、成本会加速,最优分区数落在两者之间的一段区间。

我常用的反推思路

目标分区数 ≈ max( 目标吞吐 / 单分区吞吐 ,  期望的消费并行度 )

例:目标 600 MB/s 写入,单分区实测可稳定写 ~50 MB/s
   → 至少 600 / 50 = 12 个分区
   若下游要 20 个消费者并行 → 取 max(12, 20) = 20 个分区
   再留一定余量(分区只增不减,宁可略多)

经验法则:按”目标吞吐 / 单分区吞吐”和”期望并行度”两条线取较大值,并适度留余量(因为分区数只能加不能减)。但别无脑堆到成千上万——KRaft 之后单集群能撑的分区数虽大幅提升,可选举/恢复/延迟的边际成本依然真实存在。

8. 存储机制:日志分段、保留与压缩

“读不删除”不代表无限堆积。分区日志在物理上被切成段文件(log segment),每段配有偏移量索引与时间索引;保留与清理都以段为单位进行。清理策略由 cleanup.policy 决定:

策略行为适用
delete(默认)retention.ms(时间)或 retention.bytes(大小)删除过期的整段日志、埋点、事件流等”过期即可丢”
compact(日志压缩)每个 key 只保留最新一条 value,旧值被回收需要”最新状态”的场景,如 CDC(Change Data Capture,变更数据捕获)、KV 变更流
compact,delete两者叠加既要保最新态又要设上限

deletecompact 的心智模型完全不同:前者是”按时间/大小丢老数据”,后者是”按 key 去重留最新”。内部 topic __consumer_offsets 就是典型的 compact topic——它只关心每个 (group, topic, partition) 的最新提交位点。理解这一点,也就明白为什么”消费进度”能被廉价、可靠地存在 Kafka 自己身上。

9. 交付语义:别把内部 EOS 当端到端 EOS

“消息会不会丢 / 会不会重”取决于交付语义(delivery semantics),Kafka 支持三档:

语义含义怎么达成
At-most-once(至多一次)可能丢,绝不重先提交 offset 再处理;或 acks=0
At-least-once(至少一次,默认)绝不丢,可能重先处理再提交 offset;acks=all + 重试
Exactly-once(精确一次,EOS)不丢不重幂等生产者 enable.idempotence=true + 事务 transactional.id(读-处理-写在一个事务内)

我认为最关键的洞见是:Kafka 的 exactly-once 是”Kafka→Kafka”闭环内的语义(典型是 Kafka Streams 的”消费→计算→产出”),它保证的是消息在 Kafka 内部不重不丢。一旦处理副作用落到外部系统(写库、调下游、扣款),端到端的”精确一次”就不再由 Kafka 单独保证——仍需要消费端幂等(用数据库唯一键、Redis SET、去重表等)来兜底。

因此生产里的主流不是硬上 EOS,而是 at-least-once + 消费端幂等:允许重复投递,靠业务幂等消化重复。acks 是可靠性总开关(0 最快最易丢 → 1 均衡 → all 最稳)。

10. 元数据管理:从 ZooKeeper 到 KRaft

历史上 Kafka 重度依赖 ZooKeeper 存元数据:broker 注册、topic/分区信息、Controller 选举、ISR 变更等都落在 ZK。这带来两个问题:多维护一套分布式系统,且元数据规模受 ZK 制约(分区数上限、故障恢复时元数据加载慢)。

于是社区用 KRaft(Kafka Raft) 把元数据管理内置进 Kafka 自己——由一组 controller 组成 quorum,通过 Raft 达成共识,不再需要外部 ZooKeeper。演进时间线:

版本里程碑依据
2.8(2021)KRaft 首次引入(early access,尝鲜)KIP-500
3.3(2022)KRaft 生产可用(production-ready)KIP-833
3.5ZooKeeper 模式标记为废弃(deprecated)KIP-833
4.0(2025)彻底移除 ZooKeeper,仅支持 KRaft4.0 Release

KRaft 不只是”少一个依赖”:把元数据变成一个内部的、可增量同步的日志后,单集群可支撑的分区规模上升约一个数量级,Controller 故障切换时的元数据加载也从”秒级”降到”接近毫秒级”(官方在 KIP-500 与 3.3 发布说明中给出的量级,具体数字随集群规模而变)——这也是前面 §7 说”能撑的分区数大幅提升”的底层原因。概念(Controller、选举、元数据)都还在,只是实现搬了家。

11. 设计哲学:Kafka 的几个核心取舍

把前面的点收束成一张”设计取舍”表——这是理解 Kafka”为什么长这样”的钥匙,也是我判断其设计深度的核心依据:

取舍Kafka 的选择换来了什么代价
broker 智能 vs 客户端智能dumb broker + smart client:broker 只顺序存日志,位点/重平衡/流控交给客户端broker 极简、可线性扩展、超高吞吐客户端更重、语义更复杂
队列为中心 vs 日志为中心log-centric:一份可重放日志被多方订阅多消费者、重放、当事件存储用不适合逐条 ack、复杂路由
随机结构 vs 顺序日志顺序追加 + 页缓存 + 零拷贝磁盘也能跑出内存级吞吐只能追加、不支持随机改
全局有序 vs 分区并行以分区换并行,只保证分区内有序水平扩展、高并发全局有序要牺牲并行
即时删除 vs 保留窗口消费不删、按策略保留多方独立消费 + 可重放占存储、要管保留/压缩
push vs pullpull 模型天然流控、批量、可回放无数据时需 long-poll 等待

一句话串起来:Kafka 的所有”快”和”稳”,几乎都来自”把日志顺序化、把状态客户端化、把并行分区化”这三个朴素但极致的选择。 抓住这条主线,就能推导出它几乎所有行为。

12. 生产反模式与踩坑

理解了机制,就能预判坑。以下是我归纳的最高频几类:

这些坑几乎无一例外地指向同一根源:误把”分区内有序 + 客户端管进度 + 日志按策略保留”这套机制当成了别的东西。回到 §1 的日志主线,多数坑都能提前预判。

13. 典型场景:为什么 AdTech 偏爱 Kafka

Kafka 官方定位是”实时数据管道 + 流处理”。落到广告系统,几乎每条关键链路都能看到它:

场景Kafka 扮演的角色
曝光/点击/转化事件流投放端把海量事件发布到 ad-events 类 topic,作为全站唯一事实来源(source of truth)
计费与对账计费系统消费事件流;出问题可重放历史事件重算——可重放日志的杀手级用途
实时特征 / 样本管道排序模型的实时特征、训练样本从事件流里流式加工(喂给 Flink / Kafka Streams)
实时预算控制 / 反作弊在事件流上开窗聚合,秒级感知超投、异常流量
系统解耦与削峰上游投放洪峰先落 Kafka,下游各系统按自己节奏消费,互不拖累

这也是为什么高并发广告服务的经典技术栈常写成 Netty + Kafka + Redis + Flink:Netty 抗接入、Kafka 抗事件洪峰并解耦、Redis 抗热点读、Flink 做实时计算。Kafka 在这条链路里是”缓冲 + 事实来源”的定盘星。

这五类场景的组件构成与架构设计(含 6 张架构图),单独展开在 《为什么 AdTech 偏爱 Kafka》

14. 与 RabbitMQ / RocketMQ 的取舍

别把 Kafka 当”更好的 RabbitMQ”,它们设计目标不同:

维度KafkaRabbitMQRocketMQ
内核模型分区化的持久日志传统 broker + 队列(AMQP)类 Kafka 的日志 + 更强业务特性
吞吐极高(百万级/秒起)
消息重放原生(按 offset)弱(取走即走)支持
复杂路由弱(topic + 分区为主)(exchange / routing key)
延迟 / 定时 / 优先级弱(需自建)(原生延迟消息等)
生态大数据 / 流计算最强通用后端消息电商 / 金融业务场景丰富

表中”吞吐 / 延迟”是默认配置下的相对量级,实际数字高度依赖硬件、批量、刷盘与可靠性配置——调优后三者差距会明显收窄,选型更应看内核模型与生态,而非单点跑分。

一句话选型:要极致吞吐、可重放、接大数据生态 → Kafka;要灵活路由 / 优先级 / 事务性业务消息 → RabbitMQ / RocketMQ。广告、日志、埋点、数仓这类”大流量 + 需重放”的场景,Kafka 几乎是默认答案。

15. 常见误解 ↔ 正解

常见误解正解
Kafka 保证消息全局有序只保证单个 partition 内有序;要有序就用 key 把消息路由到同一分区
消费一条消息,Kafka 就删一条不删。消息按保留 / 压缩策略留存,消费只推进 offset,可重放、可多方消费
消费者能立刻读到刚写入的最新一条只能读到 High Watermark 之前(已被 ISR 复制)的数据
同一 group 里加消费者就能无限提速上限是 partition 数;多出来的消费者会空闲
partition 越多越好过多会拖累选举 / 恢复 / 端到端延迟与文件句柄,要按吞吐目标折中(§7)
Kafka 默认 exactly-once默认是 at-least-once;EOS 需幂等 + 事务,且不等于端到端不重
开了事务就端到端不重不丢Kafka 内部 EOS ≠ 外部副作用不重;仍需消费端幂等
Kafka 一定要配 ZooKeeper3.3+ 用 KRaft4.0 已彻底移除 ZooKeeper
Kafka 适合做低延迟 RPC / 通用业务 MQ那是它的短板;请求-响应用 RPC,复杂路由/优先级用 RabbitMQ / RocketMQ
offset 是全局唯一的offset 是分区内的编号

16. 速查表

理解了”Kafka 是一份被切片、被复制、被多方订阅的分布式日志”,再叠加它”dumb broker + smart client”的设计哲学,就握住了推导它几乎所有行为的主轴。接下来的系列会沿这条主线继续深挖:生产者怎么高效写入Producer 分区策略与 Sticky Partitioner)、消费者怎么可靠读取Consumer 深挖)、怎么做到不丢不重可靠性与 Exactly-Once)、凭什么这么快性能内核)。


延伸阅读

本系列内部串读(Kafka 深挖系列,本篇是地基):

一手资料与优质教程:


views
Share this post on:

Previous Post
Kafka Producer 深挖:分区策略与 Sticky Partitioner 是怎么把延迟砍半的
Next Post
Redis 布隆过滤器实现:位图、RedisBloom 与缓存穿透防线