bg6.kafka
批处理(Spark 那套)的世界观是:数据是一个有界的、静止的数据集——先攒一天的日志成一个目录,再一次性算完。可现实里的数据不是这样:订单在持续下单、埋点在持续上报、设备在持续心跳。数据是持续到达的,不是一批批的。
于是你被迫做一件别扭的事:攒批。攒一小时算一次、攒一天算一次。攒得久 → 延迟高;攒得短 → 小文件多、调度开销大。批处理的模型和"数据流着来"这个事实之间,始终隔着一层。
要做实时,你需要一根管道:数据一产生就流进来,下游谁想算、什么时候算、算几遍,各取所需。这根管道就是本篇的主角。
传统消息队列(Message Queue,本篇简称 MQ):生产者把消息投递到 broker(服务端),消费者取走消息;取走(确认消费)后消息通常就从队列里删除。这是"队列"这个词的本意——进去、出来、没了。
这个模型在"一条命令投递给一个处理方"的点对点场景很好用。但拿它扛大数据管道,会撞三堵墙:
- 吞吐撞顶。传统 MQ 为了支持"随机确认、按需删除单条消息",broker 往往要维护每条消息的状态、做随机读写。数据量一大,随机 IO 就成瓶颈。
- 不可回溯。消息被消费后就删了。下游程序算错了想重算昨天的数据?对不起,数据已经不在了。批处理里"重跑一遍昨天的分区"是家常便饭,MQ 给不了。
- 多消费者难独立。同一份数据,风控要看、计费要看、大屏也要看。传统 MQ 的"队列"语义下,一条消息被一个消费者取走,别人就看不到了;即使用发布订阅(pub/sub)模式,各订阅方的进度、回溯也很难各自独立管理。
Kafka(2011 年由 LinkedIn 开源,现为 Apache 顶级项目)的答案是:别把它当队列,当成一个只能在末尾追加、能按位置随便回读的日志文件。
核心抽象:分布式追加日志(distributed commit log,业界也叫 commit log / append-only log / 追加写日志)。 消息不是"取走即删",而是顺序追加到日志末尾并按保留期驻留;每个消费者用一个**位置(offset,偏移量)**记录自己读到哪了,读到哪、读几遍、从头再读,都是消费者自己的事。
这一个换模型,三堵墙同时塌了(后面"为什么"一节给因果链):追加写 → 顺序 IO → 高吞吐;消息不删 → 可回溯;进度由消费者自管 → 多方独立消费互不干扰。
代价:
有序性被限制在分区内。为了并行,日志被切成多个分区,全局有序不再免费。
精确一次(exactly-once)语义复杂且有开销。默认拿到的是"至少一次",要精确一次得上幂等 + 事务,有协调成本。
低延迟点对点场景未必划算。如果你只是"A 发一条指令给 B"、要求毫秒级点对点、没有回溯和多消费需求,Kafka 的分区/副本/批量机制反而是额外负担——这种场景传统轻量 MQ 更合适(何时不该用,见边界一节)。
broker(代理节点):一个 Kafka 服务进程 / 节点。多个 broker 组成集群。
topic(主题):一类消息的逻辑名字,比如
orders、user_click。生产者往 topic 发,消费者从 topic 读。partition(分区):topic 在物理上被切成若干分区,每个分区就是一个独立的追加日志。分区是并行、扩展、有序的基本单位(承接 02 篇的分区概念)。
offset(偏移量):消息在某个分区内的位置编号,从 0 单调递增。它只在分区内有意义。
producer(生产者):发消息的一方。可以指定 key,key 决定消息落到哪个分区。
consumer(消费者):读消息的一方。
consumer group(消费者组):一组协作的消费者,共同消费一个 topic;组内每个分区只被组内一个消费者读(组内是负载均衡 / 竞争消费),不同组各自独立读全量(跨组是发布订阅)。
flowchart LR
PA["Producer A<br/>(key=user1)"] --> P0
PB["Producer B<br/>(无 key,轮询)"] --> P1
PB --> P2
subgraph T["Topic: orders(3 个分区)"]
P0["Partition 0<br/>append-only log"]
P1["Partition 1<br/>append-only log"]
P2["Partition 2<br/>append-only log"]
end
P0 --> B0
P1 --> B1
P2 --> B1
subgraph G1[".......Group: billing(2 个消费者,组内分摊)"]
B0["Consumer b0"]
B1["Consumer b1"]
end
P0 --> A0
P1 --> A0
P2 --> A0
subgraph G2["..........Group: analytics(1 个消费者,独收全量)"]
A0["Consumer a0"]
end
orders 有 3 个分区;billing 组用 2 个消费者分摊这 3 个分区(b1 扛 2 个),analytics 组用 1 个消费者独收全部 3 个分区。两个组各读各的、进度互不干扰——这正是传统 MQ 做不到的"多方独立消费 + 各自回溯"。
因果链 1:追加写 → 顺序 IO → 高吞吐
结论:Kafka 的写入是只在分区日志末尾追加,从不改历史、不随机插入。
为什么快(对照你的 OS 概念):
- 机械盘(HDD)的随机写要经历寻道 + 旋转延迟,每次几毫秒;顺序写则是磁头连续划过,省掉寻道。二者差距是数量级的——顺序写常在百 MB/s 量级,小块随机写可能掉到个位数 MB/s。(具体数字随盘型差异很大,此处只说数量级。)
- SSD 没有机械寻道,但顺序写依然占优:能合并成大块、对齐擦除块、减少写放大。
- 追加写让 Kafka 天然是顺序 IO,于是用普通磁盘也能扛很高的写吞吐,不必靠昂贵的随机 IO 优化。
代价:日志只能追加,意味着"删除 / 更新单条消息"这种操作 Kafka 不擅长——它靠**按时间/大小保留(retention)和日志压缩(log compaction)**来回收空间,而不是逐条删(见边界一节)。
因果链 2:offset 由消费者自管 → 可回溯 + 多组独立
结论:消息读没读、读到哪,不由 broker 记账删除,而是消费者维护自己的 offset(提交到内部 topic __consumer_offsets)。
为什么能回溯 / 多组独立:
- offset 只是分区里的一个整数位置。想重算昨天?把 offset 重置到昨天那条,从头再读一遍即可(
--from-beginning或 seek 到指定 offset)。数据还在日志里,没被消费"吃掉"。 - 每个消费者组维护自己的一套 offset。billing 组读到 offset=1000,不影响 analytics 组还停在 offset=200。互不干扰。
用你熟的概念近似类比(近似):offset 像你顺序读一个大文件时的文件指针 position——文件内容不因为你读过就消失,你可以 seek 回任意位置重读。Kafka 把"每个消费者一个独立文件指针"这件事做成了一等公民。
代价:offset 要持久化、要在消费者宕机/重启时正确恢复,还要处理"先提交 offset 还是先处理消息"——这个顺序问题直接决定了投递语义是"最多一次"还是"至少一次"(见底层机制)。
因果链 3:分区 → 水平扩展 + 并行
结论:一个 topic 切成 N 个分区,分散到多个 broker 上(承接 02 篇分区)。
为什么能扩展 / 并行:
- 写:不同分区在不同 broker,生产者按 key 或轮询把流量摊开,写吞吐随分区数近似线性增长。
- 读:一个消费者组里,分区是分配给消费者的最小单位。3 个分区最多让 3 个消费者并行读;想读得更快,加分区 + 加消费者。
- 分区还是有序的边界:Kafka 只保证单个分区内有序,同一个 key 的消息会落到同一分区(默认按 key 哈希),从而"同 key 有序"。
代价:并行的代价是牺牲全局有序——跨分区不保证顺序(这是最容易踩的坑,见边界一节)。分区数也不是越多越好(同上)。 一个分区可以有多个副本(replica),分布在不同 broker 上,其中一个是 leader(主副本),其余是 follower(从副本)。
- 所有读写都走 leader;follower 只干一件事:不断从 leader 拉取新消息,把自己追平。
- ISR(in-sync replicas,同步副本集):那些"跟得够紧"的副本集合(含 leader)。判定标准是 follower 落后 leader 的时间没超过
replica.lag.time.max.ms(默认 30000 ms = 30s)。落后太多的 follower 会被踢出 ISR,追上了再加回来。 - 高水位(high watermark,HW):ISR 里所有副本都已复制到的最小 offset。消费者只能读到高水位以内的消息——没被 ISR 集体确认的消息,读者看不到,避免读到"随时可能因主副本挂掉而消失"的数据。
flowchart LR
Prod["Producer<br/>acks=all"] -->|1. 写入| L["Leader<br/>(Partition 0)"]
L -->|2. 复制| F1["Follower 1<br/>(ISR)"]
L -->|2. 复制| F2["Follower 2<br/>(ISR)"]
F1 -->|3. ack| L
F2 -->|3. ack| L
L -->|4. 全部 ISR 确认后<br/>回 ack| Prod
这套 leader/follower + ISR 就是"多副本一致性"在 Kafka 里的具体落地:用副本冗余对抗单机故障,用 ISR + 高水位保证"读者看到的都是已被多数派安全持有的数据"。(集群元数据本身的一致性由 KRaft 负责)
先把"快"的机制补全(除了第三节的顺序写、offset、分区,还有两个):
- page cache + 不强制 fsync:Kafka 写入其实是写进 OS 的 page cache(页缓存),由 OS 决定何时刷盘,Kafka 默认不为每条消息强制
fsync。少了同步刷盘的等待 → 写延迟低、吞吐高。热数据读取也直接命中 page cache,不回磁盘。 - 零拷贝(zero-copy):把日志从磁盘发给消费者时,Kafka 用
sendfile系统调用,让数据在内核里直接从 page cache 送到网卡,跳过"内核 → 用户态 → 再回内核"的多次拷贝,也省去反序列化。 - 批量 + 压缩:生产者把多条消息攒成一个 batch 再发,摊薄每条消息的网络/协议固定开销;还能整批压缩。
你可能会问:不 fsync,机器一断电,page cache 里没落盘的数据不就丢了? 会。所以 Kafka 的"不丢"不靠单机刷盘,靠副本——这是它和传统数据库耐久性思路的关键分野。
"不丢"的因果链,落在生产者的 acks 参数上:
acks 取值 |
leader 何时回 ack | 后果 |
|---|---|---|
0 |
发出去就算成功,不等确认 | 最快,最可能丢 |
1 |
leader 写入本地成功即回 | leader 落盘前宕机会丢 |
all(即 -1) |
ISR 里所有副本都确认后才回 | 最稳,延迟最高 |
"又快又不丢"的完整链条(快与不丢在此收敛):
acks=all:只要 producer 收到成功 ack,消息已被 ISR 里所有副本持有。- 配
min.insync.replicas=2(broker/topic 配置,默认1):要求 ISR 至少有 2 个副本在线,否则生产者写入直接被拒(报错),宁可写不进,也不接受"只剩一份还假装成功"。 unclean.leader.election.enable=false(0.11 起默认false):leader 挂了,只从 ISR 里选新 leader,绝不让一个落后的、不在 ISR 的副本上位——那会丢掉它没跟上的数据。
代价(务必记住):
acks=all要等副本确认,延迟上升;ISR 里最慢的那个副本拖累整体。min.insync.replicas=2时,如果副本挂到只剩 1 个 ISR,生产者会被拒绝写入——这是拿可用性换持久化(对应 02 篇 CAP 的取舍:一致/持久 vs 可用,你必须二选一)。- 副本本身要占存储和网络:
replication.factor=3就是 3 倍存储、3 倍复制流量。
一句话总结这套取舍:Kafka 把耐久性从"单机刷盘"搬到了"多机副本",于是既能享受不 fsync 的快,又能靠副本兜住不丢——代价是延迟、存储、以及故障时对可用性的牺牲。
从 Kafka 3.0 起,生产者默认 enable.idempotence=true(幂等生产者),相应地 acks 默认值也从旧版的 1 变成了 all(3.0 之前默认是 acks=1)。也就是说,新版本开箱就是较强的一档。
rebalance(再均衡):消费者组内"分区 → 消费者"的分配关系发生重排。
- 何时触发:有消费者加入 / 离开(正常关闭、宕机、心跳超时被踢)、订阅的分区数变化、订阅关系变化。
- 由谁协调:broker 端的 group coordinator(组协调器)。
- 两种风格:
- eager(急切式,早期默认):stop-the-world——所有消费者先放弃全部分区,再统一重新分配,期间整组暂停消费。
- cooperative incremental(协作增量式,
CooperativeStickyAssignor):只迁移真正需要变动的分区,其余照常消费,停顿大幅减小。生产上推荐。
rebalance 风暴(一个高频坑):如果消费者一批消息处理太久,超过
max.poll.interval.ms(默认 300000 ms = 5 min),coordinator 认为它"卡死"把它踢出组 → 触发 rebalance → 它下次 poll 又想加回来 → 再次 rebalance……反复横跳,全组吞吐雪崩。相关的还有session.timeout.ms(默认 45000 ms = 45s,3.0 起)控制心跳超时。(这两个默认值为版本敏感项。)
演进方向:新的消费者组协议(KIP-848,把再均衡逻辑更多放到 broker 端、减少客户端停顿)在近几个版本逐步推进。这块仍在演进,具体某版本的可用性/默认状态请以官方文档为准.
下面处理的语义:读消息 + 计算 + 写
最多一次(at-most-once):先提交 offset,再处理。提交完就崩了、消息没处理 → 丢,但绝不重复。适合"丢一点无所谓、绝不能重"的场景(如某些计数、监控采样)。
至少一次(at-least-once):先处理,再提交 offset(生产端失败会重试)。处理完还没提交就崩 → 重启后会重读这条 → 可能重复处理,但不丢。这是最常见的默认工程实践。
精确一次(exactly-once,EOS / exactly-once semantics):不丢也不重。它由两块拼成:
- 幂等生产者(idempotent producer):每个生产者有 PID + 每分区递增序列号,broker 靠它给因重试导致的重复写去重。3.0 起默认开启。它解决的是"生产者重发造成分区内重复"。
- 事务(transactions,0.11 引入):让"读 Kafka → 处理 → 写 Kafka"(read-process-write)这一串跨分区操作原子提交——offset 提交和结果写入要么一起成功、要么一起回滚。 边界与代价(关键):
EOS 主要在 Kafka 生态内部闭环才成立(如 Kafka Streams 的
processing.guarantee=exactly_once_v2,或 07 篇 Flink 通过两阶段提交 +KafkaSink的EXACTLY_ONCE)。一旦下游是外部系统(写 MySQL、发 HTTP、调第三方),Kafka 的事务管不到它,端到端精确一次得靠下游自己幂等(如按业务主键 upsert)或下游支持事务。
代价:事务协调器开销、延迟上升、复杂度上升。不是所有场景都值得——很多实时数仓链路,"至少一次 + 下游按主键幂等"就是更划算的工程选择。
分区在物理上不是一个大文件,而是一串段文件(segment):写满一个(log.segment.bytes 默认约 1 GiB)就开新段,老段按保留期(log.retention.hours 默认 168 h = 7 天)删除或压缩。offset 就是穿过这些段、单调递增的位置编号。
一张 ASCII 图看两个消费者组在同一份日志上各自的位置:
Partition 0 的追加日志(append-only log)——offset 从左到右单调递增:
offset: 0 1 2 3 4 5 6 [LEO=7 下一条写这里]
消息: [m0] [m1] [m2] [m3] [m4] [m5] [m6]
^ ^
| |
Group billing 已提交 offset=4 Group analytics 已提交 offset=6
(下一条读 m4) (下一条读 m6)
<----- 已消费、仍驻留、可重放 ----->
(billing 想重算,把 offset seek 回 0 即可,数据都还在)
高水位(HW):ISR 全部复制到的位置,消费者只能读到 HW 以内。
LEO(Log End Offset):日志末尾,下一条消息将写入的位置。
要点:同一份日志、两个组、两个独立指针。数据没被"消费掉",只要还没过保留期,任何组都能 seek 回去重读。这就是"可回溯 + 多组独立"在物理层的样子。
下面用 Kafka 自带命令行工具演示(Kafka 4.0 起集群为 KRaft 模式,不再依赖 ZooKeeper;命令本身与旧版基本一致)。命令假设本地 broker 在 localhost:9092。
# ① 建一个 3 分区、3 副本的 topic
# --partitions 3 并行度上限=3
# --replication-factor 3 每分区 3 副本,容忍 broker 故障
kafka-topics.sh --create --topic orders \
--partitions 3 --replication-factor 3 \
--bootstrap-server localhost:9092
# ② 生产 3 条消息(每行一条,回车即发)
kafka-console-producer.sh --topic orders \
--bootstrap-server localhost:9092
> {"id":1,"amt":100}
> {"id":2,"amt":200}
> {"id":3,"amt":300}
# (Ctrl-C 退出)
# ③ analytics 组从头消费(--from-beginning 把 offset 重置到 0)
kafka-console-consumer.sh --topic orders \
--group analytics --from-beginning \
--bootstrap-server localhost:9092
模拟运行结果(三条都读到,注意跨分区顺序不保证与发送顺序一致):
{"id":1,"amt":100}
{"id":3,"amt":300}
{"id":2,"amt":200}
一句话解释:analytics 组拿到了全量 3 条。此刻若再起一个 --group billing --from-beginning,它同样能读到全部 3 条——因为消息没被删,两个组各有各的 offset。这正是"多组独立消费 + 可回溯"。至于输出顺序为何不是 1、2、3?因为三条落在不同分区,跨分区无全局有序
# --describe 打印每个分区的当前位点/末尾/滞后
kafka-consumer-groups.sh --describe --group billing \
--bootstrap-server localhost:9092
模拟运行结果:
GROUP TOPIC PARTITION CURRENT-OFFSET LOG-END-OFFSET LAG CONSUMER-ID
billing orders 0 10021 10025 4 consumer-1-...
billing orders 1 9800 15200 5400 consumer-2-...
billing orders 2 10050 10050 0 consumer-3-...
一句话解释:LAG = LOG-END-OFFSET − CURRENT-OFFSET,即"还差多少条没追上"。分区 1 的 LAG=5400,说明这个消费者处理速度跟不上生产速度——要么加消费者/加分区扩并行,要么优化处理逻辑,否则数据会越积越多。监控 lag 是 Kafka 运维的第一指标。
// Java 配置
Properties props = new Properties();
props.put("bootstrap.servers", "localhost:9092");
props.put("acks", "all"); // 等 ISR 全部确认才算成功 —— 不丢的关键
props.put("enable.idempotence", "true"); // 幂等:重试不会造成分区内重复(3.0起默认true)
props.put("retries", Integer.MAX_VALUE); // 失败重试(有幂等兜底,重试不怕重复)
// 配合 broker/topic 端:min.insync.replicas=2, unclean.leader.election.enable=false
一句话解释:acks=all + enable.idempotence=true + broker 端 min.insync.replicas=2,三者合起来才是完整的"不丢"配方——单靠 acks=all 而 min.insync.replicas=1 时,副本掉到只剩 leader 仍会"成功",隐患犹在。
分区数:不是越多越好
- 太少:并行度受限(组内消费者数 > 分区数时,多出来的消费者空闲干等),吞吐上不去。
- 太多:每个分区都是文件句柄 + 内存 + 一套副本复制连接;分区数暴涨会拖慢 leader 选举、增大元数据压力、拉高端到端延迟、让 rebalance 更慢。
- 分区只能增、不能减;而且增加分区会打乱"key → 分区"的哈希映射,导致同一个 key 的历史消息和新消息落到不同分区,破坏同 key 有序。所以要按目标吞吐预估并留余量,别指望以后随意调。
跨分区没有全局有序(最常见的认知错误)
Kafka 只保证单分区内有序。想要"同一个用户的事件严格有序"?给消息带上 key=userId,让它们落到同一分区即可(同 key 有序)。想要"整个 topic 全局有序"?只有一个办法:topic 只用 1 个分区——但那等于放弃并行,吞吐被锁死。全局有序和水平扩展在 Kafka 里是互斥的,先想清楚你到底需不需要全局有序。
rebalance 风暴
见 上文。根因多是消费者单批处理太久超过 max.poll.interval.ms。对策:减小 max.poll.records(默认 500)单批拿少点、把重活异步化、用 CooperativeStickyAssignor 减小停顿。
保留期与"数据悄悄没了"
Kafka 不是数据库,默认按时间/大小滚动删除旧数据(默认保留 7 天)。你的消费者如果停机超过保留期再来读,中间的数据已经被删,直接丢。做离线补数、灾备回放,要么调大保留期、要么下游落一份持久存储,别把 Kafka 当永久归档。
何时不该用 Kafka
- 低延迟、点对点的请求响应:一条命令投给一个服务、要毫秒级、无回溯无多消费——用轻量 MQ 或直接 RPC,Kafka 的分区/副本/批量是额外负担。
- 需要按内容随机删除/更新单条消息:追加日志模型不擅长,这是关系库/文档库的活。
- 消息量极小、又想要复杂路由(按 header 灵活路由到不同队列):传统 MQ(如 RabbitMQ)的 exchange 路由更顺手。
回到开头那根裂缝:数据流着来,需要一根管道。Kafka 就是这根管道,而真正对流做计算的是下游的流引擎——这正是 Flink。
它们的分工是:
- Kafka:可靠、可回溯、高吞吐地存住并分发事件流(source of truth / 数据源)。
- Flink:从 Kafka 持续消费、做有状态计算(窗口、聚合、join),再把结果写回 Kafka 或数仓。
为什么这对搭档天生互补,几个点在本篇已埋好伏笔:offset 可重置 → Flink 出错能从上次的位点重放(配合 checkpoint 做故障恢复);分区 → Flink 的并行度可以对齐 Kafka 分区数;事务/幂等 → Flink 通过两阶段提交 + Kafka 事务实现端到端精确一次。这些流侧细节(checkpoint、watermark、状态后端)留到 后续 展开,本篇只需记住:Kafka 负责"把流稳稳地存着、供人回放",Flink 负责"把流算出来"。
小结
- Kafka 的一句话本质:把消息系统建模成"分布式追加日志"——顺序追加、按位置回读、按期驻留。
- 快:顺序 IO + page cache(不强制 fsync)+ 零拷贝 + 批量 + 分区并行。
- 不丢:多副本 + ISR +
acks=all+min.insync.replicas≥2+ 禁止脏选主;耐久性来自副本,不是单机刷盘。 - 代价:延迟↑、存储/带宽↑、故障时可用性可能被牺牲(CAP 取舍);全局有序与并行互斥;精确一次复杂且有开销。
- 定位:它是流计算的数据源与缓冲带,可回溯、多组独立消费是它相对传统 MQ 的分水岭;但低延迟点对点、随机增删、复杂路由场景,别硬套 Kafka。