Skip to content

Kafka 深度解析 ​

#组件 · #Kafka · #消息队列 · #流处理 · #发布订阅

从存储模型到高吞吐设计,从 Producer/Consumer 到可靠性保证,全面理解 Kafka 的核心思想。

一、架构概览 ​

mermaid
flowchart TB
    subgraph Cluster["Kafka 集群"]
        subgraph B1["Broker 1"]
            P0["Topic-A Partition 0<br/>(Leader)"]
            R1["Topic-A Replica 1"]
            R2["Topic-A Replica 2"]
        end
        subgraph B2["Broker 2"]
            P1["Topic-A Partition 1<br/>(Leader)"]
            R0a["Topic-A Replica 0"]
            R2a["Topic-A Replica 2"]
        end
        subgraph B3["Broker 3"]
            P2["Topic-A Partition 2<br/>(Leader)"]
            R0b["Topic-A Replica 0"]
            R1b["Topic-A Replica 1"]
        end
    end

    ZK["ZooKeeper / KRaft<br/>元数据管理<br/>(Controller选举、分区分配)"]
    ZK -->|"协调"| Cluster

    Producer["Producer<br/>生产者"] -->|"写入消息"| Cluster
    Cluster -->|"拉取消费"| Consumer["Consumer Group<br/>消费者组<br/>每个分区仅一个消费者"]

二、存储模型(Kafka 高性能的核心) ​

2.1 分区与日志段 ​

Topic: user-events (3 分区, 2 副本)

Partition 0 的文件布局:
/var/lib/kafka/data/user-events-0/
├── 00000000000000000000.log      ← 日志段 (存储消息)
├── 00000000000000000000.index    ← 稀疏索引 (offset → 文件位置)
├── 00000000000000000000.timeindex← 时间索引
├── 00000000000000102456.log      ← 下一个日志段
├── 00000000000000102456.index
└── leader-epoch-checkpoint

2.2 顺序写入 + Page Cache ​

Producer → Broker:

1. 追加消息到 .log 文件末尾 (顺序写, 仅追加)
2. 不立即 fsync, 依赖 OS Page Cache
3. 消费者读取时大概率命中 Page Cache (零磁盘 I/O)

为什么顺序写这么快?

对比随机写顺序写
磁盘寻道每次都要磁头移动几乎没有
机械硬盘100 IOPS6000+ IOPS
SSD也比顺序慢很多极快

2.3 零拷贝 (Zero-Copy) ​

传统方式 (4 次拷贝 4 次上下文切换):
磁盘 → Read Buffer → 用户态 Buffer → Socket Buffer → 网卡

Kafka sendfile (2 次拷贝 2 次切换):
磁盘 → Read Buffer ────────────────→ Socket Buffer → 网卡
                    (DMA 直接传输)

节省: 用户态拷贝 + 2 次上下文切换

2.4 日志清理策略 ​

策略参数行为
删除 (默认)cleanup.policy=delete按时间(retention.ms)或大小(retention.bytes)删除旧段
压缩cleanup.policy=compact同 Key 只保留最新 Value (适合 CDC/快照)
两者cleanup.policy=compact,delete同时生效

日志压缩原理:

压缩前:
Key: A → v1    压缩后:
Key: B → v1    Key: A → v3
Key: A → v2 → Key: B → v2
Key: B → v2    Key: C → v1
Key: A → v3

三、生产者 ​

3.1 发送流程 ​

Producer 端:
┌──────────┐   ┌──────────┐   ┌──────────┐
│序列化     │→ │分区器     │→ │消息累加器  │→  Sender 线程异步发送
│Serializer│   │Partitioner│   │RecordAccum│   (批量, 按 Broker 分组)
└──────────┘   └──────────┘   └──────────┘
                                   │
                             ┌─────┴─────┐
                             │ BufferPool │  ← 内存池复用 ByteBuffer
                             └───────────┘

3.2 分区策略 ​

java
// 默认: 指定 key → hash(key) % partition_count
//       未指定 → sticky (2.4+) 随机选分区攒批

// 自定义分区器
public class MyPartitioner implements Partitioner {
    public int partition(String topic, Object key, byte[] keyBytes,
                         Object value, byte[] valueBytes, Cluster cluster) {
        // 自定义逻辑
    }
}

3.3 ACK 与可靠性 ​

acks行为可靠性性能
0不等任何确认,发完即返回最低最高
1Leader 写入即返回,不等副本中等高
-1 (all)所有 ISR 副本都写入才返回最高较低

3.4 幂等与事务 ​

幂等生产者(enable.idempotence=true):

每条消息带 (ProducerID, SequenceNumber)
Broker 去重: 同一个 ProducerID 的 SeqNum 只能递增

事务(transactional.id):

java
producer.initTransactions();
producer.beginTransaction();
producer.send(...);
producer.send(...);
producer.commitTransaction();  // 或 abortTransaction()
// 支持跨分区跨 Topic 原子性 (EOS: Exactly Once Semantics)

四、消费者 ​

4.1 消费者组与 Rebalance ​

Consumer Group: app-log-processor

Partition: P0  P1  P2  P3  P4  P5
           │   │   │   │   │   │
Consumer:  C1  C1  C2  C2  C3  C3

规则:

  • 一个分区只能被组内一个消费者消费
  • N 个分区最多 N 个消费者(多余的闲置)
  • 消费者数量变化/超时 → 触发 Rebalance

4.2 Rebalance 协议演进 ​

版本策略特点
Eager (旧)Stop-The-World所有消费者停止,重新分配,短暂不可用
Cooperative (2.4+)增量分配只迁移必要的分区,不停止所有消费者

4.3 Offset 管理 ​

__consumer_offsets (内部 Topic, 默认 50 分区)

存储: (GroupID, Topic, Partition) → Offset

自动提交 (enable.auto.commit=true):
  按 auto.commit.interval.ms (默认 5s)

手动提交:
  consumer.commitSync()    ← 同步, 阻塞, 重试
  consumer.commitAsync()   ← 异步, 不阻塞, 不重试

4.4 消费语义 ​

语义实现方式适用场景
At Most Once先提交 offset 再处理可丢数据(监控日志)
At Least Once (默认)先处理再提交 offset可能重复(需幂等)
Exactly Once事务 + 幂等严格一致性

五、副本与 ISR ​

5.1 ISR 机制 ​

Partition P0:
  Leader   (Broker 1) ← Producer 只写 Leader
  Replica  (Broker 2) ← ISR (In-Sync Replica)
  Replica  (Broker 3) ← ISR

ISR: 与 Leader 保持同步的副本集合
     replica.lag.time.max.ms (默认 30s) 内追上即保持

判断依据:不是按消息条数落后,而是按时间。只要副本在 replica.lag.time.max.ms 内追上了 Leader,就在 ISR 中。这避免了瞬时写入高峰导致副本被频繁踢出。

5.2 Leader 选举 ​

优先副本选举 (Preferred Replica):
  每个分区的第一个副本是 "Preferred Leader"
  auto.leader.rebalance.enable=true 时自动均衡

Unclean Leader 选举:
  unclean.leader.election.enable=false (默认)
  ISR 为空时 → 宁可不可用也不选落后副本
  unclean.leader.election.enable=true → 可用性优先,可能丢数据

5.3 LEO vs HW ​

Leader:      [m1][m2][m3][m4][m5]  LEO=5
Follower A:  [m1][m2][m3]          LEO=3
Follower B:  [m1][m2][m3][m4]      LEO=4

HW (High Watermark) = min(LEO of all ISR) = 3
→ 消费者最多能读到 m3 (HW 之前的消息已全部副本确认)

六、Consumer Group Rebalance ​

6.1 Rebalance 触发条件 ​

mermaid
sequenceDiagram
    participant C1 as Consumer 1
    participant C2 as Consumer 2
    participant G as GroupCoordinator<br/>(某Broker)
    participant T as Topic Partitions

    Note over C1,C2: 正常消费中

    C1--xC1: 崩溃!
    Note over G: heartbeat 超时<br/>session.timeout.ms

    G->>G: 触发 Rebalance
    G->>C2: 撤销当前分区分配
    G->>C2: 重新分配: P0,P1,P2 全给 C2

触发条件:

  • Consumer 加入或离开 Group
  • max.poll.interval.ms 超时(两次 poll 间隔过长)
  • session.timeout.ms 心跳超时(Consumer 假死/真死)
  • Topic 分区数变化

6.2 Eager vs Cooperative Rebalance ​

特性Eager (老)Cooperative (新, 2.4+)
行为Stop-The-World:所有 Consumer 先释放全部分区增量分配:只撤销需要重新分配的分区
影响短暂消费中断零停机,未变分区继续消费
配置默认partition.assignment.strategy=CooperativeStickyAssignor

6.3 分区分配策略 ​

策略分配方式适用
RangeAssignor按 Topic 逐个分配,分区号连续分配给同一 Consumer默认,可能导致倾斜
RoundRobinAssignor轮询分配所有分区各 Consumer 分区数均匀
StickyAssignor尽量保持上次分配结果 + 均衡Rebalance 时减少分区移动
CooperativeStickyAssignorSticky + 增量 Rebalance✅ 生产推荐

七、Exactly-Once 语义 ​

7.1 三层语义 ​

mermaid
flowchart LR
    subgraph Semantics["消息传递语义"]
        A["At Most Once<br/>最多一次<br/>可能丢消息"]
        B["At Least Once<br/>至少一次<br/>可能重复"]
        C["Exactly Once<br/>精确一次<br/>不丢不重"]
    end

    A -->|"acks=0, 不重试"| Producer
    B -->|"acks=all, 重试"| Producer
    C -->|"幂等+事务"| Producer

7.2 幂等 Producer(Idempotent) ​

java
// 开启幂等:Producer 自动去重
props.put("enable.idempotence", true);
// 内部机制:
// - 每个 Producer 分配 PID (Producer ID)
// - 每条消息带 sequence number(递增)
// - Broker 检查 (PID, TopicPartition, SeqNo) 是否重复
// - 自动设置 acks=all, max.in.flight=5, retries=MAX_INT

7.3 事务(Transactions) ​

跨分区原子写入 + consume-transform-produce 精确一次:

java
// 初始化事务 Producer
props.put("transactional.id", "tx-app-1");
producer.initTransactions();

// 事务性写入
producer.beginTransaction();
try {
    producer.send(new ProducerRecord<>("topic-A", key, valueA));
    producer.send(new ProducerRecord<>("topic-B", key, valueB));
    producer.commitTransaction();
} catch (Exception e) {
    producer.abortTransaction();
}

// Consume-Transform-Produce 模式:
// consumer 读取 → 处理 → producer 事务写入 → 提交消费 offset
producer.sendOffsetsToTransaction(offsets, consumerGroup);

八、性能参数调优 ​

6.1 Broker 端 ​

参数说明建议
num.network.threads网络处理线程CPU 核数
num.io.threads磁盘 I/O 线程CPU 核数 × 2
num.partitions总分区数单 Broker ≤ 4000
log.segment.bytes日志段大小1GB
log.retention.hours保留时间按需求 (168=7天)
socket.send.buffer.bytes发送缓冲区1MB
socket.receive.buffer.bytes接收缓冲区1MB

6.2 Producer 端 ​

参数说明建议
batch.size一批消息大小16KB-1MB
linger.ms等待攒批时间5-100ms
buffer.memory总缓冲区32MB-128MB
compression.type压缩算法lz4 或 snappy
max.in.flight.requests.per.connection最大在途请求≤ 5(开启幂等时 ≤ 5)

6.3 Consumer 端 ​

参数说明建议
fetch.min.bytes最少拉取字节1MB
fetch.max.wait.ms等待攒够 fetch.min.bytes 的最长时间500ms
max.partition.fetch.bytes每分区最大拉取10MB
max.poll.records单次 poll 最多条数500-2000
max.poll.interval.ms两次 poll 最大间隔(超时触发 rebalance)5min

九、常见问题 ​

7.1 消息丢失场景 ​

场景原因解决方案
Produceracks=0 不等待acks=-1 + min.insync.replicas=2
BrokerLeader 宕机 + unclean 选举unclean.leader.election.enable=false
Consumer先提交 offset 后处理先处理再提交(At Least Once)

7.2 消息重复与幂等 ​

重复消费不可避免 (At Least Once), 需业务幂等:

1. 唯一键去重: 每条消息带 UUID/业务 ID, DB 唯一约束
2. 版本号: Redis SETNX / DB version 字段
3. 事务标记: 处理前记录 (txID, 状态), 避免重复处理

7.3 消息积压 ​

原因处理
消费速度跟不上增加消费者(不超过分区数)
消息倾斜调整分区策略使 key 更均匀
下游慢/异常熔断降级,临时扩容
消费代码问题排查耗时逻辑,优化处理速度

7.4 消息顺序保证 ​

Kafka 只保证单个分区内消息有序:

需要全局有序 → 1 个分区 (牺牲并行度)
需要按用户有序 → 按 user_id 分区

十、Kafka vs 其他 MQ ​

维度KafkaRabbitMQRocketMQPulsar
吞吐量极高 (百万/秒)中等 (万/秒)高 (十万/秒)极高
延迟毫秒级微秒级毫秒级毫秒级
存储磁盘持久化, 长期内存+磁盘, 消费即删磁盘持久化分层存储
消费模型拉取 (Pull)推送 (Push)拉取 (Pull)拉取 (Pull)
可靠性高 (副本+ISR)极高高 (同步刷盘)高
顺序分区内有序队列内有序队列内有序分区内有序
事务✅ (0.11+)❌✅❌
适用场景日志、流、大数据管道业务消息、RPC电商、金融云原生、多租户
多租户❌ (需额外隔离)部分(vhost)❌✅ 原生

九、Kafka 生态 ​

组件作用
Kafka Connect数据导入/导出连接器(DB → Kafka → ES/HDFS)
Kafka Streams轻量级流处理库(Java,无需独立集群)
KSQL (ksqlDB)SQL 流处理(SELECT ... FROM stream WHERE ...)
Schema RegistryAvro/Protobuf Schema 管理
MirrorMaker 2跨集群数据复制
KRaft (3.3+ 生产可用)去 ZooKeeper,元数据由 Kafka 自身 Raft 协议管理

十、Kafka vs RocketMQ vs Pulsar — 三大消息队列对比 ​

三家都是"分布式消息队列",但设计哲学大不相同。选型不是在选"谁更好",而是在选"谁更适合你的场景":

维度KafkaRocketMQPulsar
底层存储日志段文件 (Segment)CommitLog + ConsumeQueueBookKeeper (Ledger)
消息模型分区 (Partition)队列 (MessageQueue)Topic → Segment
消费模型拉 (Pull) 为主拉+推 (Pull+Push)拉 (Pull)
延迟消息❌ 原生不支持✅ 18 个延迟等级✅ 任意延迟
事务消息✅ Exactly-Once (幂等+事务)✅ 半消息回查机制✅
顺序消息🟡 分区内有序✅ 局部有序✅
流量控制仅消费者生产者+消费者生产者+消费者
多租户❌ 弱🟡 通过命名空间✅ 原生 Namespace
存储与计算分离❌ 绑定❌ 绑定✅ Broker 无状态
协议兼容自有协议自有协议Kafka/Pulsar/AMQP/RabbitMQ
运维复杂度🟡 中等🟡 中等❌ 高(多组件)
社区/生态✅✅ 最大✅ 国内强🟡 追赶中
典型用户LinkedIn, Uber, Netflix阿里, 滴滴, 美团Yahoo, Verizon, Tencent

ISR 机制与数据可靠性 ​

ISR(In-Sync Replicas)是 Kafka 保证数据可靠性的核心机制,但它的行为受到多个参数的影响:

mermaid
flowchart TB
    Leader["Partition Leader<br/>接收所有写入"]

    subgraph ISR["ISR (In-Sync Replicas)"]
        F1["Follower 1<br/>拉取延迟 < replica.lag.time.max.ms(30s)"]
        F2["Follower 2<br/>拉取延迟 < 30s"]
    end

    subgraph OSR["OSR (Out-of-Sync Replicas)"]
        F3["Follower 3<br/>拉取延迟 > 30s ❌ 被踢出 ISR"]
    end

    Leader --> F1
    Leader --> F2
    Leader -.->|"延迟过大→踢出"| F3

可靠性配置速查:

场景: 绝对不能丢消息
  acks = -1 (all)
  min.insync.replicas = 2 (至少 2 个 ISR 确认)
  unclean.leader.election.enable = false (非 ISR 副本不能选为 Leader)

场景: 追求吞吐,偶尔丢可忍
  acks = 1 (Leader 确认即可)
  compression.type = lz4 (压缩减轻网络 I/O)

场景: 极低延迟 (如实时推荐)
  acks = 0 (不等确认)
  linger.ms = 0 (立即发送,不等批次)

ISR 的核心权衡:min.insync.replicas 越大越可靠,但可用性越低——如果 ISR 数量降到此值以下,写入会直接拒绝。一般生产环境设 replication.factor=3, min.insync.replicas=2,这样可以容忍 1 个节点故障。

参考 ​

批注模式

💬 文章评论

暂无评论,来说点什么吧 👇

编程学习笔记