全文共 5,051 字 预计阅读 15 分钟
bg

bg6.kafka

批处理(Spark 那套)的世界观是:数据是一个有界的、静止的数据集——先攒一天的日志成一个目录,再一次性算完。可现实里的数据不是这样:订单在持续下单、埋点在持续上报、设备在持续心跳。数据是持续到达的,不是一批批的。

于是你被迫做一件别扭的事:攒批。攒一小时算一次、攒一天算一次。攒得久 → 延迟高;攒得短 → 小文件多、调度开销大。批处理的模型和"数据流着来"这个事实之间,始终隔着一层。

要做实时,你需要一根管道:数据一产生就流进来,下游谁想算、什么时候算、算几遍,各取所需。这根管道就是本篇的主角。

传统消息队列(Message Queue,本篇简称 MQ):生产者把消息投递到 broker(服务端),消费者取走消息;取走(确认消费)后消息通常就从队列里删除。这是"队列"这个词的本意——进去、出来、没了。

这个模型在"一条命令投递给一个处理方"的点对点场景很好用。但拿它扛大数据管道,会撞三堵墙:

  1. 吞吐撞顶。传统 MQ 为了支持"随机确认、按需删除单条消息",broker 往往要维护每条消息的状态、做随机读写。数据量一大,随机 IO 就成瓶颈。
  2. 不可回溯。消息被消费后就删了。下游程序算错了想重算昨天的数据?对不起,数据已经不在了。批处理里"重跑一遍昨天的分区"是家常便饭,MQ 给不了。
  3. 多消费者难独立。同一份数据,风控要看、计费要看、大屏也要看。传统 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(主题):一类消息的逻辑名字,比如 ordersuser_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 里所有副本都确认后才回 最稳,延迟最高

"又快又不丢"的完整链条(快与不丢在此收敛):

  1. acks=all:只要 producer 收到成功 ack,消息已被 ISR 里所有副本持有。
  2. min.insync.replicas=2(broker/topic 配置,默认 1):要求 ISR 至少有 2 个副本在线,否则生产者写入直接被拒(报错),宁可写不进,也不接受"只剩一份还假装成功"
  3. 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):不丢也不重。它由两块拼成:

    1. 幂等生产者(idempotent producer):每个生产者有 PID + 每分区递增序列号,broker 靠它给因重试导致的重复写去重。3.0 起默认开启。它解决的是"生产者重发造成分区内重复"。
    2. 事务(transactions,0.11 引入):让"读 Kafka → 处理 → 写 Kafka"(read-process-write)这一串跨分区操作原子提交——offset 提交和结果写入要么一起成功、要么一起回滚。 边界与代价(关键)
  • EOS 主要在 Kafka 生态内部闭环才成立(如 Kafka Streams 的 processing.guarantee=exactly_once_v2,或 07 篇 Flink 通过两阶段提交 + KafkaSinkEXACTLY_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=allmin.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。
Back to Blog