一句话结论

Kafka 是分布式日志系统——消息按序追加到 Partition 末尾(顺序写),消费者按 Offset 拉取(Pull 模型)。高吞吐靠顺序写 + 零拷贝 + 批量压缩。

核心概念

概念

说明

Topic

消息类别(如 device_logs)

Partition

分区,每个 Partition 内有序

Producer

生产者,写消息到 Partition

Consumer Group

消费者组,组内分摊 Partition

Offset

消息在 Partition 中的位置

Broker

Kafka 节点

消息可靠性

Producer → Broker:
  acks=0: 不等确认(可能丢)
  acks=1: Leader 确认即可(Leader 宕机可能丢)
  acks=all: 所有 ISR 确认(最安全,默认 min.insync.replicas≥2)

Broker 存储:
  Partition 多副本 → ISR 机制

Consumer:
  手动提交 Offset → 处理完再提交

高吞吐原理

顺序写

Kafka 的消息是追加到 Partition 文件末尾(append-only),
是顺序写入而非随机写入——磁盘顺序写速度可达 600MB/s+,
媲美内存随机写速度。数据不删不改,天然适合日志语义。

Page Cache + 零拷贝

Producer → Broker: 消息写入 Page Cache → 异步刷盘
Consumer ← Broker: sendfile() 零拷贝 → 数据从 Page Cache 直接到网卡
                             不经过用户态内存拷贝!

批量 + 压缩

Producer 批量发送(batch.size + linger.ms)
Broker 批量存储(Segment 文件)
Consumer 批量拉取(fetch.min.bytes)
消息支持 gzip/snappy/lz4/zstd 压缩

Partition 分配策略

Consumer Group 内的 Partition 分配:
  Partition 1 ──→ Consumer A
  Partition 2 ──→ Consumer B  
  Partition 3 ──→ Consumer C
  Partition 4 ──→ Consumer A  (一个 Consumer 可消费多个 Partition)

规则: 同组内一个 Partition 只能被一个 Consumer 消费
     一个 Consumer 可以消费多个 Partition
     Consumer 数量 > Partition 数量 → 多余的 Consumer 空闲

ISR(In-Sync Replicas)

保持与 Leader 同步的副本集合。副本落后超过 replica.lag.time.max.ms(默认 30s)被踢出 ISR。Leader 宕机 → Controller 从 ISR 中选新 Leader。

ISR 机制:
  Partition 0:  Leader(Broker1)  Follower(Broker2)  Follower(Broker3)
                [0,1,2,3,4,5]    [0,1,2,3,4,5]      [0,1,2,3] ← 落后了!
                                                      ↑ 超过阈值踢出 ISR

  acks=all + min.insync.replicas=2:
    至少 Leader + 1 个 Follower 确认才算写入成功

Rebalance 详细过程

触发条件:
  - Consumer 加入或离开 Consumer Group
  - Topic Partition 增加
  - Consumer 心跳超时(session.timeout.ms 默认 45s)

Rebalance 流程(Eager 协议,Kafka 2.4+ 支持 Cooperative 增量重分配):
  1. Group Coordinator 收到触发事件
  2. 所有 Consumer 停止消费("Stop the World")
  3. Coordinator 根据 partition.assignment.strategy 重新分配
  4. Consumer 收到新分配,恢复消费
  
Eager vs Cooperative:
  Eager:   全停 → 重新分配 → 全恢复(有停顿)
  Cooperative: 只调整有变化的 Partition,其他继续消费(无停顿)

Offset 管理

Offset 提交:
  自动提交(enable.auto.commit=true,默认 5s 间隔):
    ✅ 简单
    ❌ 可能丢消息(poll 后自动提交但还没处理完)
    ❌ 可能重复消费(处理完但还没到下一个提交间隔就崩溃)

  手动提交:
    ✅ 精确控制(处理完再提交)
    ❌ 代码复杂
    ⚠️ 注意: 提交的是下次要消费的 offset(不是已消费的最后一条)

__consumer_offsets Topic:
  Kafka 内部用这个 Topic 存 Consumer Group 的 Offset
  (Kafka 0.9 之前 offset 存 Zookeeper → 性能瓶颈)

消息语义

语义

实现方式

难度

At most once

auto commit + 处理前提交

简单

At least once

手动 commit + 处理后提交

中等

Exactly once

幂等 Producer + 事务 + 消费端去重

复杂

幂等 Producer(enable.idempotence=true):
  每条消息带 ProducerID + SequenceNumber
  Broker 去重(同一 PID 的同一 SN 只持久化一次)

事务 Producer:
  beginTransaction → 发送多条消息 → commitTransaction
  原子性保证(所有消息要么都可见要么都不可见)

与 RabbitMQ 对比

Kafka

RabbitMQ

模型

分布式日志(append-only)

消息队列(消费后删除)

消费模式

Pull(消费者按需拉取)

Push(Broker 推送)

消息回溯

✅ 按 offset/时间回溯

❌ 消息确认后删除

吞吐量

百万级 msg/s

万级 msg/s

路由灵活性

Topic → Partition(简单)

Exchange → Binding → Queue(灵活)

单消息 ACK

❌ 只能按 offset 确认

✅ 单条 ACK/NACK

延时队列

❌ 需自行实现

✅ TTL+DLX 原生支持

典型场景

日志采集、流处理、事件溯源

业务任务、RPC 调用、延迟任务

高频面试问题

Q: Kafka 为什么吞吐这么高?

30 秒回答: 四个核心设计:① 顺序写(磁盘顺序写接近内存速度)② Page Cache + sendfile 零拷贝(数据从磁盘到网卡不经过用户态)③ 批量压缩(Producer/Consumer/Broker 三端批量)④ Partition 并行(水平扩展无上限)。

Q: Kafka 的 ISR 和 Raft 的 Quorum 有什么区别?

30 秒回答: ISR 是动态集合(跟不上就踢出),Raft Quorum 是固定多数派(2N+1 个节点过半数确认)。ISR 追求低延迟(等最快的几个),Raft 追求安全性(等多数派里的多数)。

Q: Rebalance 期间能消费吗?

30 秒回答: Eager Rebalance(老方案):不能,全部暂停。Cooperative Rebalance(Kafka 2.4+):可以,只暂停被重新分配的那部分 Partition。

项目相关

目前项目用 RabbitMQ。如果未来做设备日志采集、Agent 事件溯源会考虑 Kafka。

速记

Kafka=分布式日志+顺序写+零拷贝。Partition=分区有序。Consumer Group=分摊消费。ISR=同步副本。acks=all最安全。Rebalance=重新分配暂停消费。vs RabbitMQ=业务任务用RMQ,日志流用Kafka。