在广告系统里,一次曝光(impression)、一次点击(click)、一次转化(conversion)都是一个事件。我关心的问题是:一天几十上百亿条这样的事件,如何可靠地、按序地、能被多方同时消费地从投放端流向计费、报表、风控、特征管道与实时数仓。Apache Kafka 最初在 LinkedIn 诞生,正是为了统一处理这类高吞吐、需持久化、要被多方订阅与重放的事件流——广告事件流只是它最典型的应用场景之一。
本文以 Jay Kreps 提出的「日志为中心」抽象为线索,把 Kafka 从一条 append-only 日志推演到分布式流平台,并沿途标出我认为真正决定生产行为的机制与取舍。
TL;DR
- 内核是一个”分布式、可持久化、可重放的提交日志(commit log)“。抓住”append-only 的有序日志”,其余概念都是它的自然推论;它的设计哲学是 “dumb broker + smart client”——broker 只顺序存日志,位点/重平衡/流控全下放给客户端,这是它能扛超高吞吐的根本原因。
- Topic 是逻辑消息流,Partition 是分布式的最小单位——并行、扩展、顺序、容量都由它决定且互相牵制;其上 Offset / ISR / High Watermark 三件套定义了”消费者能看到什么、什么算不丢”:offset 是位点,**ISR(In-Sync Replicas)**是选主候选池,HW(High Watermark)决定消费者可见边界。
- Kafka 是 pull(拉)模型:消费者自己控制速率与回放,换来天然的批量、流控与重放能力。
- 交付语义默认 at-least-once(至少一次,绝不丢但可能重);端到端不丢不重需”幂等/事务生产 + 幂等消费”,别把 Kafka 内部的 EOS(Exactly-Once Semantics) 当成端到端 EOS。
- 分区数是绕不开的容量决策:太少限制并行,太多拖累选举/恢复/延迟/文件句柄——要按目标吞吐反推。
- 元数据管理已从 ZooKeeper 迁到 KRaft(Kafka 内置、基于 Raft 的元数据模式):2.8 引入、3.3 生产可用、4.0 起彻底移除 ZooKeeper。
- 动手比读更快:§6 有一份可直接复制运行的 5 分钟 demo(KRaft 单节点 + CLI +
acks=all幂等 Producer + 先处理后提交的 Consumer)。
Table of contents
Open Table of contents
- 1. 内核:Kafka 是一个”分布式提交日志”
- 2. 从”日志”到”分布式事件流平台”
- 3. 核心概念全景
- 4. 消息模型:为什么 Kafka 选了发布-订阅
- 5. 消费者组与再平衡(Rebalance)
- 6. 上手:把概念跑起来(KRaft + CLI + Java)
- 7. 分区数:一个绕不开的容量决策
- 8. 存储机制:日志分段、保留与压缩
- 9. 交付语义:别把内部 EOS 当端到端 EOS
- 10. 元数据管理:从 ZooKeeper 到 KRaft
- 11. 设计哲学:Kafka 的几个核心取舍
- 12. 生产反模式与踩坑
- 13. 典型场景:为什么 AdTech 偏爱 Kafka
- 14. 与 RabbitMQ / RocketMQ 的取舍
- 15. 常见误解 ↔ 正解
- 16. 速查表
- 延伸阅读
1. 内核:Kafka 是一个”分布式提交日志”
抛开所有花哨的名词,我把 Kafka 的心脏归结为一个极简的数据结构:一条只追加(append-only)、按写入顺序排列、每条记录带一个单调递增编号(offset)的日志。
就这么一个”账本”,配合三个关键决定,就长成了 Kafka:
- 写入只追加到末尾 —— 顺序写磁盘,快到接近写内存。
- 读不删除、按策略留存 —— 消费不是”取走”,而是”读到某个位置”。同一份日志因此能被 N 个互不相干的消费者各读各的。
- 把这份日志切片、复制、分散到多台机器 —— 切片得到 Partition(拿到并行与扩展),复制得到 Replica(拿到高可用)。
一句话:Kafka = 把一个”可持久化、可重放的日志”,做成了分布式、可水平扩展、多方可订阅的系统。后面所有概念,都是围着这句话展开的。
这也解答了一个经典追问——“Kafka 到底是 MQ、数据库还是流处理引擎?“我的结论是:它本质是日志,MQ / 存储 / 流处理只是这份日志的不同用法。这个”日志为中心(log-centric)“的世界观,是 Kafka 与传统”队列为中心”的 MQ 最根本的分野(§11 会展开)。
1.1 一条 Record 长什么样
日志里的每个元素叫 Record(记录/消息),它不是裸字符串,而是一个结构:
| 字段 | 含义 |
|---|---|
| Key(可空) | 决定分区路由与压缩语义的键(如 user_id、order_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 里这几类能力常常同时用:
- ① 发布/订阅:投放端把曝光事件发布到 Kafka。
- ② 持久化存储:计费与对账系统需要能重放昨天的事件重算。
- ③ 流式处理:实时预算控制、反作弊在事件流上开窗聚合。
- ④ 系统集成:再由 Connect 把结果落库到数仓/湖。
3. 核心概念全景
先用一张总览图建立全局空间感——上半部分是各核心概念之间的关系与数量基数(cardinality),下半部分是生产/消费的部署视角——再逐个精确拆解:
上图 B 部分的元数据管理画的是 KRaft(Kafka 3.3+ 生产可用、4.0 起彻底移除 ZooKeeper,见 §10);更早的版本这里是一个外部的 ZooKeeper 集群负责注册与选主,但除此之外的所有关系(Topic / Partition / Replica / Consumer Group 等)完全不变。
3.1 Topic(主题)
消息的逻辑分类 / 命名频道。生产者往 topic 发、消费者从 topic 订阅,类比”一个具名的事件流”(如 ad-impressions、ad-clicks)。Topic 是逻辑概念,真正存数据的是它下面的 Partition。
3.2 Partition(分区)—— 分布式的最小单位
一个 topic 被切成一个或多个 partition,每个 partition 就是第 1 节那条 append-only 有序日志。 它是 Kafka 分布式、并行、扩展、顺序、容量的共同基本单位,而这些属性彼此牵制:
- 并行:不同 partition 可分布在不同 broker,读写并行展开——高吞吐的来源;也是一个 group 内消费并行度的上限。
- 扩展:加 broker、加 partition 即可水平扩展(注意 partition 数只增不减,减少要重建 topic)。
- 顺序:Kafka 只保证单个 partition 内有序,不保证跨 partition 的全局有序。
这条”只保证分区内有序”是 Kafka 最重要、也最常被误解的性质。要顺序,就得让相关消息落到同一个 partition(靠 key);要并行,就得把消息摊到多个 partition——二者天然冲突,这正是顺序性与吞吐取舍的根。分区数怎么定,见 §7。
3.3 Offset(偏移量)
每条消息在其 partition 内唯一、单调递增的编号。 三个要点:
- Offset 是分区内的,不是全局的。
- 需要区分两个位置:current position(下次要读的位置) 与 committed offset(已提交、已确认处理到的位置)——两者之差正是消费重复/丢失风险的来源。
- 消费进度由消费者提交 offset 来记录,默认存于内部 topic
__consumer_offsets(一个 compact topic,见 §8)。
正因进度归消费者自己管,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(消费者 / 消费者组)
- Consumer:从 partition 拉取(pull) 消息的客户端。
- Consumer Group:一组协作消费同一 topic 的消费者,用
group.id标识。核心规则只有一条:同一个 group 内,一个 partition 只会被一个 consumer 消费。
于是:想让多个系统各拿一份完整数据(广播)→ 用不同 group;想让一个系统内并行分担负载 → 同一 group 内加消费者(上限 = partition 数)。为什么 Kafka 选 pull 而非 push,见 §5.1。
一图看清”组内消费者数与分区数”的关系:
3.7 Replica / ISR / High Watermark —— 可见性与一致性的边界
- Replica(副本):partition 的冗余拷贝,分散在不同 broker 以容灾。
- Leader / Follower:每个 partition 有一个 leader,默认生产者和消费者都只与 leader 交互;follower 只从 leader 拉数据保持同步(2.4+ 起消费者可选择就近从 follower 读,用于跨机架优化,但写入始终走 leader)。
- ISR(In-Sync Replicas,同步副本集):与 leader 保持同步的副本集合(含 leader)。Leader 挂了只从 ISR 里选新 leader——这是”不丢消息”的前提。
- LEO(Log End Offset):日志末尾、下一条将写入的偏移。
- High Watermark(HW,高水位):已被所有 ISR 复制、因而对消费者可见的最高 offset。消费者只能读到 HW 之前的数据;HW 与 LEO 之间的尾部消息(已在 leader、还没被所有 ISR 复制)对消费者不可见。
这套机制用一张时序图最清楚:
为什么消费者只能读到 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 如何用一套机制把两者统一:
队列模型:消息进队列,一条只被一个消费者取走。适合任务分发,但想让”多个系统各拿一份完整数据”就很别扭(得为每个消费者复制一份)。
发布-订阅模型(Kafka):以 topic 为载体,发布一条,所有订阅方各自都能收到。当一个 topic 只有一个订阅组时,它退化成队列模型——所以 pub-sub 在功能上向下兼容队列模型。
Kafka 用 topic(广播)+ consumer group(组内负载均衡) 一举同时表达两种语义:不同 group = 广播(各拿全量);同 group 多消费者 = 队列式负载均衡(分摊)。
RocketMQ 的消息模型与 Kafka 基本一致;差异之一是 Kafka 没有”队列”这个概念,与之对应的是 Partition。
5. 消费者组与再平衡(Rebalance)
消费者可随时加入或退出 group。当组成员变化(或订阅的分区数变化)时,分区所有权会被重新分配给存活消费者——这个过程叫 rebalance(再平衡)。
它带来高可用与弹性伸缩,但也有代价,分两种策略:
- Eager rebalance(积极再平衡):所有消费者先停下、放弃全部分区、重新入组再领新分配——有一个短暂的全组 stop-the-world 窗口。
- Cooperative rebalance(合作 / 增量再平衡):分多阶段,只挪动一小部分分区,其余照常消费,避免全组停摆。普通消费者需显式配置
CooperativeStickyAssignor才启用(默认仍是 eager 的RangeAssignor);Kafka Streams 则已默认采用。
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、RedisSETNX、或专门的去重表(§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 端每分区一个批次,攒批效率下降 |
我常用的反推思路:
目标分区数 ≈ 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 | 两者叠加 | 既要保最新态又要设上限 |
delete与compact的心智模型完全不同:前者是”按时间/大小丢老数据”,后者是”按 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.5 | ZooKeeper 模式标记为废弃(deprecated) | KIP-833 |
| 4.0(2025) | 彻底移除 ZooKeeper,仅支持 KRaft | 4.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 pull | pull 模型 | 天然流控、批量、可回放 | 无数据时需 long-poll 等待 |
一句话串起来:Kafka 的所有”快”和”稳”,几乎都来自”把日志顺序化、把状态客户端化、把并行分区化”这三个朴素但极致的选择。 抓住这条主线,就能推导出它几乎所有行为。
12. 生产反模式与踩坑
理解了机制,就能预判坑。以下是我归纳的最高频几类:
- 热点分区(key skew):用低基数或倾斜的 key(如按国家、按大客户 ID)→ 少数分区被打爆,其余闲置。对策:选高基数、分布均匀的 key,或引入 salting(给 key 加盐/加随机后缀打散)。
- 分区数拍脑袋:上线随手设 3 个分区,后期吞吐见顶又无法缩减、只能重建 topic 迁移。对策:按 §7 反推并留余量。
- 把 Kafka 当 RPC / 低延迟请求-响应:Kafka 是高吞吐管道,不是低延迟点对点通道;同步等应答会很别扭。对策:请求-响应用 RPC,事件流才用 Kafka。
- 消费者数 > 分区数:多出的消费者永远空闲,误以为加机器能提速。对策:先加分区,再加消费者。
- 无界保留 / 无压缩:默认或过长 retention 让磁盘悄悄撑爆;状态类 topic 忘开 compaction。对策:按业务显式设
retention.*或cleanup.policy。 - 单批处理太慢触发无限 rebalance:一次
poll处理太久超过max.poll.interval.ms,被踢出组 → 重平衡 → 更慢,恶性循环。对策:减小单批量、异步处理、调大间隔。 - 把内部 EOS 当端到端不重:以为开了事务就万事大吉,外部副作用仍会重复。对策:消费端做幂等(§9)。
这些坑几乎无一例外地指向同一根源:误把”分区内有序 + 客户端管进度 + 日志按策略保留”这套机制当成了别的东西。回到 §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”,它们设计目标不同:
| 维度 | Kafka | RabbitMQ | RocketMQ |
|---|---|---|---|
| 内核模型 | 分区化的持久日志 | 传统 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 一定要配 ZooKeeper | 3.3+ 用 KRaft,4.0 已彻底移除 ZooKeeper |
| Kafka 适合做低延迟 RPC / 通用业务 MQ | 那是它的短板;请求-响应用 RPC,复杂路由/优先级用 RabbitMQ / RocketMQ |
| offset 是全局唯一的 | offset 是分区内的编号 |
16. 速查表
- 内核:一个分布式、可持久化、可重放的 append-only 有序日志。
- 设计哲学:dumb broker + smart client、log-centric、顺序换吞吐、分区换并行、保留换可重放、pull 换流控。
- Record:key(路由 + 压缩主键)/ value / timestamp / headers / offset。
- Partition:分布式最小单位,决定并行 / 扩展 / 顺序 / 容量;数量只增不减,按吞吐反推(§7)。
- Offset:分区内位置;区分 position 与 committed;消费不删数据 → 可重放、多方独立消费。
- Replica / ISR / HW:Leader 收发、ISR 是选主候选池、HW 决定消费者可见边界。
- Consumer:pull 模型;组内一 partition 只被一人消费;增减成员触发 rebalance。
- 存储:日志分段;
cleanup.policy = delete(按时间/大小)或compact(按 key 留最新)。 - 交付语义:默认 at-least-once;EOS 是 Kafka 内部语义,端到端仍需消费端幂等。
- 元数据:KRaft(3.3 生产可用,4.0 移除 ZooKeeper)。
- 选型:极致吞吐 + 可重放 + 大数据生态 → Kafka。
理解了”Kafka 是一份被切片、被复制、被多方订阅的分布式日志”,再叠加它”dumb broker + smart client”的设计哲学,就握住了推导它几乎所有行为的主轴。接下来的系列会沿这条主线继续深挖:生产者怎么高效写入(Producer 分区策略与 Sticky Partitioner)、消费者怎么可靠读取(Consumer 深挖)、怎么做到不丢不重(可靠性与 Exactly-Once)、凭什么这么快(性能内核)。
延伸阅读
本系列内部串读(Kafka 深挖系列,本篇是地基):
- Kafka Producer 深挖:分区策略与 Sticky Partitioner:生产者怎么决定”消息进哪个分区”、怎么批量高效地把消息塞进日志。
- Kafka Consumer 深挖:poll 循环、offset 提交与 rebalance:消费侧怎么可靠、高效、不抖动地把消息读出来。
- Kafka 可靠性与 Exactly-Once 深挖:
acks/ ISR /min.insync.replicas/ 幂等生产者 / 事务与 EOS 的完整机制。 - Kafka 性能内核:磁盘系统凭什么跑出内存级吞吐:顺序写、页缓存、零拷贝、日志分段四板斧。
- 为什么 AdTech 偏爱 Kafka:广告事件中枢的五大典型场景与架构设计:把上面的原理落到广告五大场景与组件架构(Netty / Kafka / Redis / Flink)。
一手资料与优质教程:
- Apache Kafka. Introduction:官方对”分布式事件流平台”、topic / partition、生产消费模型的权威定义。
- Apache Kafka. Documentation — Design:持久化、日志、复制、消费者模型(含 pull 取舍)与交付语义的官方设计章节。
- Jay Kreps. The Log: What every software engineer should know about real-time data’s unifying abstraction:Kafka”日志为中心”世界观的思想源头。
- Apache Kafka. KIP-500: Replace ZooKeeper with a Self-Managed Metadata Quorum 与 KIP-833: Mark KRaft as Production Ready:KRaft 演进的一手提案,§10 时间线的出处。
- Redpanda. Kafka Tutorial:系统的 Kafka 概念与实操教程,含分区策略、消费者分配等专题。