编程 Apache Kafka KRaft 深度实战:从 ZooKeeper 到自管理集群,Raft 共识、Controller Quorum 与生产迁移全指南

2026-07-27 13:44:19 +0800 CST views 16

Apache Kafka KRaft 深度实战:从 ZooKeeper 到自管理集群,Raft 共识、Controller Quorum 与生产迁移全指南

一、背景:Kafka 为什么需要「摘掉」ZooKeeper?

如果你用过 Kafka,一定见过这个经典架构图:Producer → Broker → Consumer,旁边挂着一个 ZooKeeper 集群。很多人习惯了这种搭配,觉得 ZooKeeper 不过是「存点元数据」,没什么大不了的。

但真实的生产环境里,ZooKeeper 恰恰是 Kafka 集群中最脆弱的环节。

1.1 ZooKeeper 的三宗罪

第一宗罪:额外运维成本

一个 Kafka 集群需要维护两套分布式系统——Kafka 本身 + ZooKeeper ensemble。ZooKeeper 需要奇数节点(通常是 3 或 5),每个节点有自己的 JVM 参数、磁盘 IO 模式、网络配置。这意味着你的运维复杂度直接翻倍。

更令人头疼的是 ZooKeeper 的「脑裂」问题。虽然 ZooKeeper 本身有 Zab 协议保证一致性,但在实际生产中,由于网络分区、磁盘延迟波动等原因,ZK 集群的 Leader 选举经常导致 Kafka Controller 的重新选举,进而引发整个集群的「雪崩」。

第二宗罪:Controller 选举的性能瓶颈

在传统架构中,Kafka Controller 的选举依赖 ZooKeeper 的临时节点(ephemeral node)。所有 Broker 在 ZK 的 /controller 路径上争抢创建节点,谁创建成功谁就是 Controller。这个过程涉及 ZK 的 Zab 共识,通常需要几百毫秒到几秒。

问题是,Controller 挂了之后,新的 Controller 需要从 ZK 读取全量元数据(所有 Topic、Partition、Replica 的分配信息),然后重建内存状态。对于拥有上万个 Partition 的生产集群,这个过程可能长达数分钟。在这段时间里,整个集群的元数据操作(创建 Topic、分区重分配、Preferred Leader 选举)全部不可用。

第三宗罪:元数据一致性的根本缺陷

Kafka 的元数据分散在两个地方:一部分在 ZooKeeper 中(Topic、Partition、ACL、Quota 等),另一部分在 Kafka Broker 的内存中。这种「分布存储」导致了一个经典的一致性问题:

假设你通过 kafka-admin.sh 创建了一个 Topic。这个操作写入 ZK 后返回成功。但是 Kafka Controller 可能还没从 ZK Watch 中感知到这个变更。如果你立即对这个 Topic 发送消息,Broker 可能会返回 UNKNOWN_TOPIC_OR_PARTITION 错误。

这种「写后读不一致」在传统架构中无法从根本上解决,因为 ZK 和 Kafka 之间没有事务性保证。

1.2 为什么要用 KRaft?

KRaft(Kafka Raft Metadata mode)是 Kafka 社区从 2.8 版本开始引入、到 3.x 版本逐步稳定、最终在 4.0 中完全取代 ZooKeeper 的元数据管理模式。

核心思路很简单:把 ZooKeeper 的职能内化到 Kafka 自身。Kafka 实现了一个基于 Raft 共识算法的元数据复制层,让一部分 Broker 节点组成 Controller Quorum,专门负责元数据的存储和分发。

这带来的好处是革命性的:

  • 运维简化:一个集群,一套配置,一个监控
  • 元数据访问延迟降低 10-100 倍(Raft 走内网 TCP,比 ZK 的 Zab 快得多)
  • Controller 选举从秒级降到毫秒级
  • 元数据变更的事务性保证(写入即可见)
  • 单节点即可运行(开发环境不再需要额外启动 ZK)

二、核心概念:Raft 共识算法与 KRaft 架构

2.1 Raft 共识算法速通

Raft 是 Diego Ongaro 在 2013 年提出的分布式共识算法,设计目标是「比 Paxos 更容易理解」。Kafka 选择 Raft 而不是继续用 Zab,主要原因是 Raft 的工程实现更简洁,且社区有大量成熟的参考实现。

Raft 的核心机制可以用三个子问题概括:

Leader Election(领导者选举)

集群中的每个节点有三种角色:

  • Leader:唯一的写入点,所有客户端请求都发给 Leader
  • Follower:被动复制 Leader 的日志
  • Candidate:选举过程中的临时角色

选举流程如下:

1. Follower 在 Election Timeout(150~300ms 随机)内没收到 Leader 心跳
2. Follower → Candidate,Term +1,给自己投票并发 RequestVote RPC
3. 收到多数派(N/2 + 1)投票的 Candidate 成为 Leader
4. Leader 开始发送 AppendEntries(心跳)维持权威

关键设计:随机超时时间。每个节点的超时时间是随机的,极大减少了「同时发起选举导致 Split Vote」的概率。

Log Replication(日志复制)

Leader 收到客户端请求后:

  1. 追加到本地日志
  2. 并行发送 AppendEntries RPC 给所有 Follower
  3. 等待多数派确认写入
  4. 日志状态变为 Committed
  5. 应用到状态机(Applied)

这里有一个很重要的概念叫 Quorum。对于 3 节点的集群,Quorum = 2(多数派);对于 5 节点,Quorum = 3。只要 Quorum 存活,集群就能正常工作。

Safety(安全性保证)

Raft 保证了以下关键特性:

  • Election Safety:每个 Term 最多一个 Leader
  • Leader Append-Only:Leader 从不覆盖或删除日志
  • Log Matching:两个日志在相同 Index/Term 的条目内容必然相同
  • Leader Completeness:被选举的 Leader 必须包含所有已 Committed 的日志

2.2 KRaft 的架构分层

Kafka 的 KRaft 实现将元数据管理抽象为三层:

┌─────────────────────────────────────┐
│        Metadata API Layer           │
│  (CreateTopic, AlterConfig, ACL)    │
├─────────────────────────────────────┤
│      Metadata Record Layer          │
│  (序列化/反序列化元数据记录)          │
├─────────────────────────────────────┤
│       Raft Log Layer                │
│  (Batch 追加、复制、快照)           │
├─────────────────────────────────────┤
│       Network Transport Layer       │
│  (基于 Kafka 自定义协议 + TCP)      │
└─────────────────────────────────────┘

Metadata API Layer

这是最上层,处理客户端(Admin Client、Broker、Kafka Tool)发来的元数据变更请求。包括:

  • CreateTopicsRequest:创建 Topic
  • AlterConfigsRequest:修改配置
  • CreateAclsRequest:创建 ACL
  • CreatePartitionsRequest:增加分区

这些请求最终会转化为 Metadata Record,写入 Raft Log。

Metadata Record Layer

每个元数据操作被序列化为一条 Metadata Record。Record 的 Schema 由 Kafka 协议定义,包含:

Record {
  type: 操作类型 (TopicRecord, PartitionRecord, ConfigRecord...)
  version: Schema 版本
  key: 操作对象的唯一标识
  value: 操作内容的 Protobuf 序列化
  timestamp: 操作时间
}

Raft Log Layer

这是 KRaft 的核心。Kafka 在 Raft 的基础上做了大量工程优化:

Batch Append:多个元数据记录被打包成一个 Batch,一次性写入 Raft Log,减少 IO 次数

Unflushed Log:Leader 在内存中维护 Recent Log,异步刷盘。Follower 也类似,先写 Page Cache 再刷盘

Periodic Snapshot:Raft Log 不能无限增长。KRaft 定期生成 Metadata Snapshot,把当前全量元数据序列化成一个文件。新的节点加入时,不需要 replay 全量 Log,只需加载最新 Snapshot + 后续增量 Log

// Kafka 源码中 MetadataSnapshot 的核心逻辑
public class MetadataSnapshot {
    private final long lastContainedLogTimestamp;
    private final Map<String, TopicMetadata> topics;
    private final Map<String, ConfigResource> configs;
    private final Map<String, AclBinding> acls;
    private final Map<Uuid, BrokerMetadata> brokers;
    
    public void serialize(Records records) {
        // 使用 Kafka 自研的 Binary Protocol 编码
        // 支持增量更新、零拷贝
    }
}

Network Transport Layer

KRaft 没有使用 gRPC 或 HTTP 作为 Raft 的传输协议,而是复用了 Kafka 自有的二进制协议。这意味着:

  • 同一端口可以同时处理数据流量和 Raft 流量(端口复用)
  • 不需要额外建立连接池
  • 可以复用 Kafka 已有的认证和加密机制(SASL、SSL)

2.3 Controller Quorum:集群的「大脑」

在 KRaft 模式下,一个很重要的设计是 Process Role(进程角色)。每个 Kafka 进程可以扮演:

  • Controller:只参与元数据管理,不处理数据读写
  • Broker:只处理数据读写,不参与元数据投票
  • Combined(Broker + Controller):同时承担两种角色
Controller Quorum(3个 Controller 节点)
│
├── Controller-1 (Leader) ← 所有元数据写入经过此节点
├── Controller-2 (Follower)
└── Controller-3 (Follower)
         │
         ▼
Broker 集群(N 个 Broker 节点)
├── Broker-1 → 从 Controller Quorum 订阅元数据变更
├── Broker-2 → 订阅
├── Broker-3 → 订阅
└── Broker-N → 订阅

这种架构的核心优势是 元数据与数据流量的完全解耦。Controller Quorum 的成员可以独立部署在专门的机器上,不受 Broker 数据流量波动的影响。

元数据分发机制

KRaft 中,元数据的传播不是通过 ZooKeeper Watch 实现的,而是通过 Metadata Fetch API。每个 Broker 定期(默认 metadata.max.age.ms = 5 分钟)或被动(收到 Metadata Update 通知)向 Controller 拉取最新的元数据映像。

// Broker 端的元数据更新逻辑
public class MetadataUpdater {
    private final MetadataCache cache;
    private final ControllerNodeProvider controllerProvider;
    
    public CompletableFuture<MetadataUpdate> poll(Optional<Long> lastSeenOffset) {
        // 1. 向 Controller 发送 MetadataFetch RPC
        // 2. 传入 lastSeenOffset(上次看到的 Log Offset)
        // 3. Controller 返回从该 Offset 开始的增量变更
        // 4. Broker 应用变更到本地 MetadataCache
    }
    
    public void onMetadataChanged() {
        // Controller 也会主动推送变更通知
        // Broker 收到后立即发起 MetadataFetch
    }
}

相比 ZooKeeper Watch 的「推模式」,KRaft 采用了 推拉结合 的策略:

  • :Controller 在关键元数据变更时主动通知 Broker
  • :Broker 定期拉取确保不会错过任何变更

2.4 节点发现机制的变化

ZooKeeper 时代,Broker 通过 ZK 的 /brokers/ids 路径发现其他节点。KRaft 模式中,节点发现通过 Metadata Log 完成:

当一个新的 Broker 启动时:

  1. 读取配置文件中的 controller.quorum.bootstrap.servers,找到至少一个 Controller
  2. 向 Controller 发送 BrokerRegistration RPC
  3. Controller 在 Metadata Log 中追加一条 BrokerRegistrationRecord
  4. 所有其他 Broker 通过日志复制感知到新节点加入

这种方式的好处是:节点注册本身就是元数据变更日志的一部分,天然具有事务性和持久性。

三、架构分析:为什么 KRaft 比 ZooKeeper 快?

这个问题从理论到工程都有清晰的答案。

3.1 协议层面的对比

特性ZooKeeper (Zab)KRaft (Raft)
共识协议Zab (ZooKeeper Atomic Broadcast)Raft
Log 结构全局顺序 WAL分 Batch 追加
选举时间秒级(依赖 ZK 自身选举)毫秒级(随机超时)
元数据复制ZK 与 Kafka 双路径单路径(Raft Log)
事务保证ZK 保证,但 Kafka 不直接使用内建于 Raft Log
运维节点ZooKeeper ensemble(3-5节点)Controller Quorum(1-3节点)
可观测性ZK 四字命令 + JMXKafka Metrics + Controller Metrics

3.2 延迟对比的工程原因

原因一:少了一跳

传统模式下,创建 Topic 的操作路径是:

Admin Client → Kafka Broker → ZooKeeper (Create) → ZK Ack → Broker → Ack

KRaft 模式下:

Admin Client → Controller (Raft Append) → Quorum Ack → Ack

KRaft 少了一层网络中转,直接由 Controller 处理元数据写入。

原因二:Batch 优化

ZooKeeper 的每个操作都是独立的,不支持批量写入。KRaft 的 Controller 可以将多个元数据操作合并成一个 Raft Batch:

// KRaft 的批处理逻辑
public class RaftBatchAppender {
    private final List<ApiMessage> pendingRecords = new ArrayList<>();
    
    public void append(ApiMessage record) {
        pendingRecords.add(record);
        if (pendingRecords.size() >= maxBatchSize || 
            timeSinceLastFlush >= maxBatchLatencyMs) {
            flush();
        }
    }
    
    private void flush() {
        // 1. 创建一个 Raft Batch
        // 2. 批量序列化所有 Record
        // 3. 一次 fsync 写入本地 Log
        // 4. 批量发送 AppendEntries RPC
        MemoryRecords batch = MemoryRecordsBuilder.build(
            pendingRecords.stream()
                .map(this::serialize)
                .collect(toList())
        );
        raftLog.append(batch);
        sendAppendEntries(batch);
        pendingRecords.clear();
    }
}

这个优化在生产环境能极大地提升吞吐量。对于大规模的元数据操作(比如批量创建 1000 个 Topic),KRaft 比 ZooKeeper 快了 50 倍以上。

原因三:本地读

ZooKeeper 中,Kafka Controller 读取元数据需要通过网络请求 ZK。KRaft 模式下,Controller Quorum 的每个节点都在本地内存中维护了完整的 Metadata Cache。读取元数据是纯内存操作,没有网络开销。

四、代码实战:从零搭建 KRaft 集群

4.1 单节点 KRaft 集群(开发环境)

首先,下载 Kafka 3.9+(包含稳定版 KRaft 支持):

wget https://downloads.apache.org/kafka/3.9.0/kafka_2.13-3.9.0.tgz
tar -xzf kafka_2.13-3.9.0.tgz
cd kafka_2.13-3.9.0

KRaft 模式需要首先生成 Cluster ID:

# 生成一个 UUID 作为集群标识
KAFKA_CLUSTER_ID=$(bin/kafka-storage.sh random-uuid)
echo $KAFKA_CLUSTER_ID
# 输出示例: MkU3OEVBNTcwNTJENDM2Qk

初始化 Log 目录:

# 使用默认配置初始化
bin/kafka-storage.sh format -t $KAFKA_CLUSTER_ID \
  -c config/kraft/server.properties

查看 config/kraft/server.properties 的核心配置:

# 进程角色:Combined 模式(既是 Controller 又是 Broker)
process.roles=broker,controller

# 节点 ID(必须唯一)
node.id=1

# Controller Quorum 配置
controller.quorum.voters=1@localhost:9093

# 数据目录
log.dirs=/tmp/kraft-combined-logs

# 监听器配置
listeners=PLAINTEXT://0.0.0.0:9092,CONTROLLER://0.0.0.0:9093
advertised.listeners=PLAINTEXT://localhost:9092

# 不同监听器的 Security Protocol 映射
listener.security.protocol.map=CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT,SSL:SSL,SASL_PLAINTEXT:SASL_PLAINTEXT,SASL_SSL:SASL_SSL

# 内部控制通道
inter.broker.listener.name=PLAINTEXT

启动:

bin/kafka-server-start.sh config/kraft/server.properties

验证启动成功:

# 创建测试 Topic
bin/kafka-topics.sh --create \
  --bootstrap-server localhost:9092 \
  --replication-factor 1 \
  --partitions 3 \
  --topic test-topic

# 查看 Topic 列表
bin/kafka-topics.sh --list \
  --bootstrap-server localhost:9092

4.2 三节点生产级集群

生产环境建议将 Controller 和 Broker 分离部署。这里演示一个 3 Controller + 3 Broker 的集群。

Controller 节点配置(controller-1.properties):

process.roles=controller
node.id=1
controller.quorum.voters=1@controller1:9093,2@controller2:9093,3@controller3:9093
log.dirs=/data/kraft/controller-logs

# Controller 只需要 CONTROLLER 监听器
listeners=CONTROLLER://0.0.0.0:9093
listener.security.protocol.map=CONTROLLER:PLAINTEXT

# 元数据日志配置
log.segment.bytes=1073741824  # 1GB
log.retention.ms=604800000    # 7天
log.cleaner.enable=true

# 快照配置
metadata.log.max.record.bytes.between.snapshots=20971520  # 20MB 变更后触发快照

Broker 节点配置(broker-1.properties):

process.roles=broker
node.id=11
controller.quorum.voters=1@controller1:9093,2@controller2:9093,3@controller3:9093
log.dirs=/data/kafka/data-01,/data/kafka/data-02  # 多数据目录

# Broker 监听器
listeners=PLAINTEXT://0.0.0.0:9092
advertised.listeners=PLAINTEXT://broker1:9092
listener.security.protocol.map=PLAINTEXT:PLAINTEXT

# 元数据同步
metadata.max.age.ms=30000  # 30秒主动拉取元数据

逐个初始化并启动:

# 在所有节点上
KAFKA_CLUSTER_ID=$(bin/kafka-storage.sh random-uuid)
# 注意:所有节点必须使用相同的 CLUSTER_ID!

# Controller 节点
bin/kafka-storage.sh format -t $CLUSTER_ID -c config/controller-1.properties
bin/kafka-server-start.sh config/controller-1.properties

# Broker 节点
bin/kafka-storage.sh format -t $CLUSTER_ID -c config/broker-1.properties
bin/kafka-server-start.sh config/broker-1.properties

4.3 使用 Docker 部署 KRaft

对于容器化环境,Kafka 官方提供了 KRaft 模式的 Docker 镜像:

# docker-compose.yml
version: '3.8'
services:
  controller-1:
    image: apache/kafka:3.9.0
    hostname: controller-1
    container_name: controller-1
    environment:
      CLUSTER_ID: 'MkU3OEVBNTcwNTJENDM2Qk'
      KAFKA_NODE_ID: 1
      KAFKA_PROCESS_ROLES: 'controller'
      KAFKA_LISTENERS: 'CONTROLLER://0.0.0.0:9093'
      KAFKA_CONTROLLER_QUORUM_VOTERS: '1@controller-1:9093,2@controller-2:9093,3@controller-3:9093'
      KAFKA_CONTROLLER_LISTENER_NAMES: 'CONTROLLER'
      KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: 'CONTROLLER:PLAINTEXT'
      KAFKA_LOG_DIRS: '/var/lib/kafka/data'
    volumes:
      - controller-1-data:/var/lib/kafka/data
    networks:
      - kafka-net

  controller-2:
    image: apache/kafka:3.9.0
    hostname: controller-2
    container_name: controller-2
    environment:
      CLUSTER_ID: 'MkU3OEVBNTcwNTJENDM2Qk'
      KAFKA_NODE_ID: 2
      KAFKA_PROCESS_ROLES: 'controller'
      KAFKA_LISTENERS: 'CONTROLLER://0.0.0.0:9093'
      KAFKA_CONTROLLER_QUORUM_VOTERS: '1@controller-1:9093,2@controller-2:9093,3@controller-3:9093'
      KAFKA_CONTROLLER_LISTENER_NAMES: 'CONTROLLER'
      KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: 'CONTROLLER:PLAINTEXT'
      KAFKA_LOG_DIRS: '/var/lib/kafka/data'
    volumes:
      - controller-2-data:/var/lib/kafka/data
    networks:
      - kafka-net

  controller-3:
    image: apache/kafka:3.9.0
    hostname: controller-3
    container_name: controller-3
    environment:
      CLUSTER_ID: 'MkU3OEVBNTcwNTJENDM2Qk'
      KAFKA_NODE_ID: 3
      KAFKA_PROCESS_ROLES: 'controller'
      KAFKA_LISTENERS: 'CONTROLLER://0.0.0.0:9093'
      KAFKA_CONTROLLER_QUORUM_VOTERS: '1@controller-1:9093,2@controller-2:9093,3@controller-3:9093'
      KAFKA_CONTROLLER_LISTENER_NAMES: 'CONTROLLER'
      KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: 'CONTROLLER:PLAINTEXT'
      KAFKA_LOG_DIRS: '/var/lib/kafka/data'
    volumes:
      - controller-3-data:/var/lib/kafka/data
    networks:
      - kafka-net

  broker-1:
    image: apache/kafka:3.9.0
    hostname: broker-1
    container_name: broker-1
    depends_on:
      - controller-1
      - controller-2
      - controller-3
    environment:
      CLUSTER_ID: 'MkU3OEVBNTcwNTJENDM2Qk'
      KAFKA_NODE_ID: 11
      KAFKA_PROCESS_ROLES: 'broker'
      KAFKA_LISTENERS: 'PLAINTEXT://0.0.0.0:9092'
      KAFKA_ADVERTISED_LISTENERS: 'PLAINTEXT://broker-1:9092'
      KAFKA_CONTROLLER_QUORUM_VOTERS: '1@controller-1:9093,2@controller-2:9093,3@controller-3:9093'
      KAFKA_CONTROLLER_LISTENER_NAMES: 'CONTROLLER'
      KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: 'CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT'
      KAFKA_INTER_BROKER_LISTENER_NAME: 'PLAINTEXT'
      KAFKA_LOG_DIRS: '/var/lib/kafka/data'
      KAFKA_METADATA_MAX_AGE_MS: '30000'
    volumes:
      - broker-1-data:/var/lib/kafka/data
    ports:
      - "9092:9092"
    networks:
      - kafka-net

  broker-2:
    image: apache/kafka:3.9.0
    hostname: broker-2
    container_name: broker-2
    depends_on:
      - controller-1
      - controller-2
      - controller-3
    environment:
      CLUSTER_ID: 'MkU3OEVBNTcwNTJENDM2Qk'
      KAFKA_NODE_ID: 12
      KAFKA_PROCESS_ROLES: 'broker'
      KAFKA_LISTENERS: 'PLAINTEXT://0.0.0.0:9092'
      KAFKA_ADVERTISED_LISTENERS: 'PLAINTEXT://broker-2:9092'
      KAFKA_CONTROLLER_QUORUM_VOTERS: '1@controller-1:9093,2@controller-2:9093,3@controller-3:9093'
      KAFKA_CONTROLLER_LISTENER_NAMES: 'CONTROLLER'
      KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: 'CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT'
      KAFKA_INTER_BROKER_LISTENER_NAME: 'PLAINTEXT'
      KAFKA_LOG_DIRS: '/var/lib/kafka/data'
      KAFKA_METADATA_MAX_AGE_MS: '30000'
    volumes:
      - broker-2-data:/var/lib/kafka/data
    networks:
      - kafka-net

  broker-3:
    image: apache/kafka:3.9.0
    hostname: broker-3
    container_name: broker-3
    depends_on:
      - controller-1
      - controller-2
      - controller-3
    environment:
      CLUSTER_ID: 'MkU3OEVBNTcwNTJENDM2Qk'
      KAFKA_NODE_ID: 13
      KAFKA_PROCESS_ROLES: 'broker'
      KAFKA_LISTENERS: 'PLAINTEXT://0.0.0.0:9092'
      KAFKA_ADVERTISED_LISTENERS: 'PLAINTEXT://broker-3:9092'
      KAFKA_CONTROLLER_QUORUM_VOTERS: '1@controller-1:9093,2@controller-2:9093,3@controller-3:9093'
      KAFKA_CONTROLLER_LISTENER_NAMES: 'CONTROLLER'
      KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: 'CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT'
      KAFKA_INTER_BROKER_LISTENER_NAME: 'PLAINTEXT'
      KAFKA_LOG_DIRS: '/var/lib/kafka/data'
      KAFKA_METADATA_MAX_AGE_MS: '30000'
    volumes:
      - broker-3-data:/var/lib/kafka/data
    networks:
      - kafka-net

volumes:
  controller-1-data:
  controller-2-data:
  controller-3-data:
  broker-1-data:
  broker-2-data:
  broker-3-data:

networks:
  kafka-net:
    driver: bridge

启动:

docker-compose up -d

先等 10 秒让 Controller Quorum 完成选举,然后验证:

# 在 broker-1 上创建 Topic
docker exec broker-1 kafka-topics.sh --bootstrap-server localhost:9092 \
  --create --topic orders --partitions 6 --replication-factor 2

# 查看 Topic 详情
docker exec broker-1 kafka-topics.sh --bootstrap-server localhost:9092 \
  --describe --topic orders

4.4 生产消息示例(Go 语言客户端)

KRaft 模式对客户端完全透明。下面展示 Go 语言的生产者/消费者代码:

package main

import (
    "context"
    "fmt"
    "log"
    "time"
    
    "github.com/segmentio/kafka-go"
)

// KafkaConfig 封装连接配置
type KafkaConfig struct {
    Brokers []string
    Topic   string
}

func NewProducer(config KafkaConfig) *kafka.Writer {
    return &kafka.Writer{
        Addr:     kafka.TCP(config.Brokers...),
        Topic:    config.Topic,
        Balancer: &kafka.RoundRobin{},
        // 关键生产配置
        BatchSize:       100,           // 批量发送
        BatchTimeout:    50 * time.Millisecond,
        Async:           false,         // 同步等待确认
        RequiredAcks:    kafka.RequireAll, // -1: 等待所有副本确认
        MaxAttempts:     3,             // 重试次数
        WriteTimeout:    10 * time.Second,
    }
}

func NewConsumer(config KafkaConfig, groupID string) *kafka.Reader {
    return kafka.NewReader(kafka.ReaderConfig{
        Brokers:   config.Brokers,
        Topic:     config.Topic,
        GroupID:   groupID,
        MinBytes:  10 * 1024,  // 10KB
        MaxBytes:  10 * 1024 * 1024, // 10MB
        // 关键消费配置
        MaxWait:       1 * time.Second,
        ReadLagInterval: 1 * time.Second,
        CommitInterval:  1 * time.Second,
        StartOffset:    kafka.FirstOffset,
        IsolationLevel: kafka.ReadCommitted, // 只读取已提交消息
    })
}

func main() {
    cfg := KafkaConfig{
        Brokers: []string{"localhost:9092"},
        Topic:   "orders",
    }

    // 启动生产者
    producer := NewProducer(cfg)
    defer producer.Close()

    // 启动消费者
    consumer := NewConsumer(cfg, "order-processor")
    defer consumer.Close()

    // 发送 1000 条订单消息
    ctx := context.Background()
    for i := 0; i < 1000; i++ {
        msg := kafka.Message{
            Key:   []byte(fmt.Sprintf("order-%d", i)),
            Value: []byte(fmt.Sprintf(`{"order_id":"ORD-%05d","amount":%.2f,"ts":%d}`, 
                i, float64(i)*19.99, time.Now().UnixMilli())),
            Headers: []kafka.Header{
                {Key: "source", Value: []byte("go-producer")},
                {Key: "version", Value: []byte("1.0")},
            },
        }
        
        err := producer.WriteMessages(ctx, msg)
        if err != nil {
            log.Printf("发送失败: %v", err)
            continue
        }
        
        if i%100 == 0 {
            log.Printf("已发送 %d 条消息", i)
        }
    }

    // 消费消息
    for i := 0; i < 100; i++ {
        msg, err := consumer.ReadMessage(ctx)
        if err != nil {
            log.Printf("消费失败: %v", err)
            break
        }
        log.Printf("收到消息: key=%s, value=%s, partition=%d, offset=%d",
            string(msg.Key), string(msg.Value), msg.Partition, msg.Offset)
    }
}

4.5 元数据状态检查

KRaft 模式提供了新的工具来查看元数据状态:

# 查看 Controller Quorum 状态
bin/kafka-metadata-quorum.sh --bootstrap-server localhost:9092 describe --status

# 输出示例:
# ClusterId: MkU3OEVBNTcwNTJENDM2Qk
# LeaderId: 1
# LeaderEpoch: 42
# HighWatermark: 12345
# MaxFollowerLag: 0
# Voters: [1, 2, 3]
# CurrentVoters: [1, 2, 3]

# 查看元数据 Log
bin/kafka-dump-log.sh --cluster-metadata-decoder \
  --files /tmp/kraft-combined-logs/__cluster_metadata-0/*.log

# 查看 Broker 注册信息
bin/kafka-broker-api-versions.sh --bootstrap-server localhost:9092

关键指标解读:

  • LeaderEpoch:Controller Leader 的任期号,每次选举递增
  • HighWatermark:已提交的 Log Offset,所有 Follower 都已确认
  • MaxFollowerLag:Follower 最大落后条数,正常应接近 0
  • Voters/CurrentVoters:配置的投票节点 / 当前在线的投票节点

五、性能优化:KRaft 模式下的调优清单

5.1 Controller 节点调优

# controller.properties 生产调优

# Raft 相关配置
raft.message.max.batch.size=1048576     # 1MB: Raft 批量大小
log.flush.interval.messages=10000       # 每 10000 条消息刷一次盘
log.flush.interval.ms=1000              # 最多 1 秒刷一次盘

# 元数据缓存
metadata.log.max.record.bytes.between.snapshots=268435456  # 256MB 变更后触发快照
metadata.log.snapshot.max.new.record.bytes=268435456

# 网络配置
controller.quorum.append.linger.ms=5     # 5ms 延迟聚合
controller.quorum.request.timeout.ms=2000  # 2s 请求超时
controller.quorum.retry.backoff.ms=100     # 100ms 重试退避

# 选举配置
controller.quorum.election.timeout.ms=1000    # 1s 选举超时
controller.quorum.fetch.timeout.ms=2000       # 2s 拉取超时

5.2 Broker 端元数据同步优化

# broker.properties

# 元数据拉取频率
metadata.max.age.ms=15000              # 15 秒主动拉取(默认 5 分钟)
metadata.max.idle.interval.ms=30000    # 30 秒空闲拉取

# 元数据缓存大小
metadata.cache.max.size.bytes=536870912  # 512MB 元数据缓存

# 网络连接池
metadata.network.request.timeout.ms=5000     # 5s 元数据 RPC 超时
metadata.network.max.wait.ms=1000            # 1s 最长等待

5.3 关键性能指标与基准测试

# 创建用于压测的 Topic
bin/kafka-topics.sh --bootstrap-server localhost:9092 \
  --create --topic benchmark --partitions 12 --replication-factor 2

# 生产者性能测试
bin/kafka-producer-perf-test.sh \
  --topic benchmark \
  --num-records 1000000 \
  --record-size 1024 \
  --throughput -1 \
  --producer-props bootstrap.servers=localhost:9092 \
  acks=all linger.ms=10 batch.size=65536

# 消费者性能测试
bin/kafka-consumer-perf-test.sh \
  --topic benchmark \
  --messages 1000000 \
  --broker-list localhost:9092

对比数据(来自 Kafka 官方 Benchmark):

场景ZooKeeper 模式KRaft 模式提升
Controller 选举时间2-5 秒200-500ms10x
批量创建 1000 个 Topic45 秒1.2 秒37x
元数据变更延迟 (p99)120ms8ms15x
单节点部署需要 ZK无需 ZKN/A

5.4 生产环境监控

KRaft 模式暴露了新的 JMX Metrics:

kafka.controller:type=KafkaController,name=ActiveControllerCount
kafka.controller:type=KafkaController,name=GlobalPartitionCount
kafka.controller:type=KafkaController,name=GlobalTopicCount
kafka.controller:type=KafkaController,name=OfflinePartitionsCount

# Raft 相关 Metrics
kafka.server:type=RaftMetrics,name=CommitTimeMs
kafka.server:type=RaftMetrics,name=AppendTimeMs
kafka.server:type=RaftMetrics,name=LogFlushTimeMs
kafka.server:type=RaftMetrics,name=ElectionTimeMs
kafka.server:type=RaftMetrics,name=FetchTimeMs

Prometheus 抓取配置示例:

# prometheus.yml
scrape_configs:
  - job_name: 'kafka-controller'
    metrics_path: '/metrics'
    static_configs:
      - targets:
        - 'controller1:8080'
        - 'controller2:8080'
        - 'controller3:8080'

Grafana 告警规则:

# 关键告警
- alert: ControllerQuorumDegraded
  expr: count(kafka_controller_quorum_active) < 3
  for: 30s
  labels: severity: critical
  annotations:
    summary: "Controller Quorum 不完整,当前在线 {{ $value }}/3"

- alert: MetadataSyncLagHigh
  expr: kafka_server_raftmetrics_max_follower_lag > 1000
  for: 1m
  labels: severity: warning
  annotations:
    summary: "元数据同步延迟过高 ({{ $value }} 条)"

六、从 ZooKeeper 到 KRaft 的迁移指南

6.1 迁移前的评估

不是所有场景都适合立即迁移。先对照这张表做评估:

条件适合迁移建议暂缓
Kafka 版本≥ 3.5< 3.0
ZooKeeper 版本3.8+3.5 以下
使用 ZK 强依赖功能依赖 ZK 直接读取元数据的自定义工具
集群规模< 1000 个 Broker大规模集群需充分测试
可用性要求允许短时间维护窗口7×24 不允许停机

6.2 迁移步骤

Kafka 从 3.5 开始提供了从 ZooKeeper 到 KRaft 的双向迁移工具(KIP-833)。核心思路是 ZooKeeper ↔ KRaft 双模式运行,逐步切换:

第一阶段:准备

# 1. 对所有 Broker 开启 ZooKeeper Migration 模式
# 在每个 broker 的 server.properties 中添加:
zookeeper.metadata.migration.enable=true

# 2. 生成 Cluster ID(KRaft 模式需要)
KAFKA_CLUSTER_ID=$(bin/kafka-storage.sh random-uuid)

# 3. Broker 需要额外配置 Controller Quorum 地址
controller.quorum.voters=1@controller1:9093,2@controller2:9093,3@controller3:9093

第二阶段:部署 Controller Quorum

启动独立的 Controller 节点,它们会从 ZooKeeper 读取现有元数据:

# 使用 migration 模式启动 Controller
bin/kafka-server-start.sh config/controller-migration.properties

其中 controller-migration.properties 的关键配置:

process.roles=controller
node.id=1
controller.quorum.voters=1@controller1:9093,2@controller2:9093,3@controller3:9093
zookeeper.connect=zookeeper1:2181,zookeeper2:2181,zookeeper3:2181
zookeeper.metadata.migration.enable=true

在 migration 模式下,Controller 同时从 ZK 和 Raft Log 读取元数据,确保两边一致。

第三阶段:滚动升级 Broker

逐个重启 Broker,关闭 ZooKeeper 连接,完全切换到 KRaft 模式:

# 注释或删除 zookeeper.connect
# zookeeper.connect=...
process.roles=broker
controller.quorum.voters=1@controller1:9093,2@controller2:9093,3@controller3:9093

第四阶段:验证

# 检查所有节点状态
bin/kafka-broker-api-versions.sh --bootstrap-server broker1:9092

# 确认 Controller Quorum
bin/kafka-metadata-quorum.sh --bootstrap-server broker1:9092 describe --status

# 检查元数据是否完整
bin/kafka-topics.sh --bootstrap-server broker1:9092 --list
bin/kafka-configs.sh --bootstrap-server broker1:9092 --describe --all

# 验证 ACL 和 Quota 是否迁移成功
bin/kafka-acls.sh --bootstrap-server broker1:9092 --list

第五阶段:关闭 ZooKeeper

确认一切正常后,安全关闭 ZooKeeper 集群。建议至少观察 72 小时再关闭 ZK。

6.3 迁移踩坑清单

我整理了一份常见的迁移问题列表:

问题 1:迁移后元数据不一致

现象:Broker 启动后报 MismatchedMetadataException

原因:部分自定义管理工具直接操作 ZK 写入元数据,绕过 Kafka API

解决:检查所有 Admin 操作是否都通过 Kafka AdminClient 完成。直接写 ZK 的操作在 KRaft 模式下不可用。

问题 2:Controller Quorum 无法建立

现象:Controller 日志反复输出 No leader elected

原因:通常是网络防火墙导致 Controller 之间的 Raft 端口不通

解决:检查 controller.quorum.voters 配置中的端口是否可互通,确保 CONTROLLER 监听器暴露在正确的 IP/端口上。

# 验证 Controller 连通性
nc -zv controller1 9093
nc -zv controller2 9093
nc -zv controller3 9093

问题 3:消费者组位移丢失

现象:迁移后部分消费者组的 Offset 丢失

原因:offsets.topic.replication.factor 配置不当导致位移 Topic 未成功复制

解决:迁移前设置更高的副本因子:

bin/kafka-configs.sh --bootstrap-server localhost:9092 \
  --entity-type topics --entity-name __consumer_offsets \
  --alter --add-config replication.factor=3

问题 4:元数据同步风暴

现象:Broker 启动时全部同时拉取元数据,导致 Controller CPU 飙升

原因:滚动重启时 Broker 同时重启,元数据变更通知全部集中

解决:控制滚动重启速度,每次只重启 1 个 Broker,间隔至少 30 秒。

七、总结与展望

7.1 KRaft 模式的现状

截至 2026 年 7 月,KRaft 已经经历过 3.x 到 4.0 的多个版本迭代,生产稳定性经过了大规模验证。Confluent、腾讯云、阿里云等主流 Kafka 服务商已经全面切换到 KRaft 模式。

如果你还在运行 ZooKeeper 模式的 Kafka,现在是时候规划迁移了。不要等到 ZK 集群出问题时才后悔——KRaft 带来的运维简化、性能提升和一致性保证,值得你投入时间做迁移。

7.2 未来发展方向

KRaft 自身也在快速进化中:

  1. Controller 自动扩缩容:允许动态增加/减少 Controller 节点,无需重启集群
  2. 元数据分级存储:热数据存内存 + 冷数据存 SSD,降低大规模集群的内存开销
  3. Raft 多流水线:将元数据变更按 Topic 分组到不同的 Raft Group,突破单 Controller 的吞吐瓶颈
  4. Serverless Kafka:依托 KRaft 的轻量级元数据层,实现按需分配的计算存储分离架构

7.3 你的下一步

如果我说了这么多,你只记住一件事,那就是:今天就用 Docker 跑一个 KRaft 模式的 Kafka。5 分钟就能体验到「没有 ZooKeeper 的 Kafka」到底有多清爽。

# 最简单的开始方式
docker run -d --name kafka \
  -e CLUSTER_ID=$(uuidgen) \
  -e KAFKA_NODE_ID=1 \
  -e KAFKA_PROCESS_ROLES='broker,controller' \
  -e KAFKA_CONTROLLER_QUORUM_VOTERS='1@localhost:9093' \
  -e KAFKA_LISTENERS='PLAINTEXT://0.0.0.0:9092,CONTROLLER://0.0.0.0:9093' \
  -e KAFKA_ADVERTISED_LISTENERS='PLAINTEXT://localhost:9092' \
  -e KAFKA_LISTENER_SECURITY_PROTOCOL_MAP='CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT' \
  -e KAFKA_INTER_BROKER_LISTENER_NAME='PLAINTEXT' \
  -e KAFKA_LOG_DIRS='/var/lib/kafka/data' \
  -p 9092:9092 \
  apache/kafka:latest

# 验证一下
kafka-topics.sh --bootstrap-server localhost:9092 --list

就是这么简单。去掉 ZooKeeper 的 Kafka,就像去掉脚镣的舞者——轻盈、稳定、强大。

推荐文章

Vue3中的v-bind指令有什么新特性?
2024-11-18 14:58:47 +0800 CST
如何使用go-redis库与Redis数据库
2024-11-17 04:52:02 +0800 CST
PHP openssl 生成公私钥匙
2024-11-17 05:00:37 +0800 CST
windon安装beego框架记录
2024-11-19 09:55:33 +0800 CST
Go的父子类的简单使用
2024-11-18 14:56:32 +0800 CST
PHP中获取某个月份的天数
2024-11-18 11:28:47 +0800 CST
npm速度过慢的解决办法
2024-11-19 10:10:39 +0800 CST
程序员茄子在线接单