欢迎光临
我们一直在努力

Kafka入门集群搭建-学习指南

这篇文章主要是总结个人实际在搭建kafka集群时遇到的问题以及整体的搭建流程,借助查阅各方资料以及借助ai的帮助,最终完成kafka简单集群的搭建,文章是借助ai编写大体内容之后个人进行了细节上的修改之后完成的,希望能对各位有所帮助。


Docker Compose 搭建 Kafka KRaft 3 节点集群

本文将带你用 Docker Compose 快速搭建一个基于 KRaft 协议的 Kafka 3 节点集群,并附上一些常见的命令和问题的排查。


一、前言:为什么选择 KRaft?

自 Kafka 3.3 起,KRaft(Kafka Raft Metadata Mode)已生产可用,它用内置的 Raft 共识协议彻底取代了 ZooKeeper。这意味着:

  • 🚀 架构简化:无需单独维护 ZK 集群
  • 📈 性能提升:元数据操作更快,支持百万级分区
  • 🔧 运维友好:单一组件,配置和监控更统一

这里是以 Confluent 官方镜像 confluentinc/cp-kafka:7.5.0 为例,在VMware虚拟机上搭建 Combined 角色(同时作为 Broker 和 Controller)的 3 节点集群,适合新手学习。


二、环境准备

2.1 环境要求

  • VMware16.0 + FinalShell3.8.3 +Windows11
  • Docker Engine:20.10.11+
  • Docker Compose:V2(推荐)或 V1

至于磁盘空间与运行空间,个人设置的是虚拟机空间40gb,运行内存4gb。

检查版本:

docker version
docker compose version # 若未安装,按下方操作
# 或
docker-compose –version

安装 Docker Compose(如果缺失)

# 安装 V2 插件(推荐,命令为 docker compose)
sudo yum install -y docker-compose-plugin

# 或安装 V1 独立版(命令为 docker-compose)
sudo curl -L "https://github.com/docker/compose/releases/download/1.29.2/docker-compose-$(uname -s)$(uname -m)" \\
-o /usr/local/bin/docker-compose
sudo chmod +x /usr/local/bin/docker-compose


三、生成合法的 CLUSTER_ID

KRaft 要求集群 ID 必须是 Base64 编码的 UUID(如 dQOUz7b6Rx6kqFJqXf9y3w),不能使用标准十六进制格式。

# 方法一:用 Kafka 镜像自带工具
docker run –rm confluentinc/cp-kafka:7.5.0 kafka-storage random-uuid
# 记下输出,例如:dQOUz7b6Rx6kqFJqXf9y3w

# 方法二:用 Python3(如果环境支持)
python3 -c "import uuid, base64; print(base64.urlsafe_b64encode(uuid.uuid4().bytes).decode().rstrip('='))"


四、编写 docker-compose.yml

4.1 创建项目目录并编辑文件

mkdir ~/kafka-kraft && cd ~/kafka-kraft
vi docker-compose.yml

4.2 完整配置文件

将下面的内容粘贴进去,注意替换三处 CLUSTER_ID 为你刚生成的 Base64 UUID。

version: '3'
services:
kafka1:
image: confluentinc/cpkafka:7.5.0
container_name: kafka1
ports:
"9092:9092"
environment:
KAFKA_NODE_ID: 1
KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: 'CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT'
KAFKA_ADVERTISED_LISTENERS: 'PLAINTEXT://kafka1:9092' # 必须用服务名,不能用 localhost
KAFKA_LISTENERS: 'PLAINTEXT://0.0.0.0:9092,CONTROLLER://0.0.0.0:9093'
KAFKA_PROCESS_ROLES: 'broker,controller'
KAFKA_CONTROLLER_QUORUM_VOTERS: '1@kafka1:9093,2@kafka2:9093,3@kafka3:9093'
KAFKA_CONTROLLER_LISTENER_NAMES: 'CONTROLLER'
KAFKA_INTER_BROKER_LISTENER_NAME: 'PLAINTEXT'
KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 3
KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 3
KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 2
CLUSTER_ID: '你的Base64-UUID' # 必须替换
volumes:
./data/kafka1:/var/lib/kafka/data
restart: unlessstopped

kafka2:
image: confluentinc/cpkafka:7.5.0
container_name: kafka2
ports:
"9094:9094" # 宿主机端口映射
environment:
KAFKA_NODE_ID: 2
KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: 'CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT'
KAFKA_ADVERTISED_LISTENERS: 'PLAINTEXT://kafka2:9094' # 注意容器内监听 9094
KAFKA_LISTENERS: 'PLAINTEXT://0.0.0.0:9094,CONTROLLER://0.0.0.0:9093'
KAFKA_PROCESS_ROLES: 'broker,controller'
KAFKA_CONTROLLER_QUORUM_VOTERS: '1@kafka1:9093,2@kafka2:9093,3@kafka3:9093'
KAFKA_CONTROLLER_LISTENER_NAMES: 'CONTROLLER'
KAFKA_INTER_BROKER_LISTENER_NAME: 'PLAINTEXT'
KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 3
KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 3
KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 2
CLUSTER_ID: '你的Base64-UUID'
volumes:
./data/kafka2:/var/lib/kafka/data
restart: unlessstopped

kafka3:
image: confluentinc/cpkafka:7.5.0
container_name: kafka3
ports:
"9095:9095"
environment:
KAFKA_NODE_ID: 3
KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: 'CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT'
KAFKA_ADVERTISED_LISTENERS: 'PLAINTEXT://kafka3:9095' # 容器内监听 9095
KAFKA_LISTENERS: 'PLAINTEXT://0.0.0.0:9095,CONTROLLER://0.0.0.0:9093'
KAFKA_PROCESS_ROLES: 'broker,controller'
KAFKA_CONTROLLER_QUORUM_VOTERS: '1@kafka1:9093,2@kafka2:9093,3@kafka3:9093'
KAFKA_CONTROLLER_LISTENER_NAMES: 'CONTROLLER'
KAFKA_INTER_BROKER_LISTENER_NAME: 'PLAINTEXT'
KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 3
KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 3
KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 2
CLUSTER_ID: '你的Base64-UUID'
volumes:
./data/kafka3:/var/lib/kafka/data
restart: unlessstopped

配置要点:

  • KAFKA_PROCESS_ROLES: 'broker,controller' → Combined 节点
  • CONTROLLER://0.0.0.0:9093 → KRaft 内部通信端口
  • PLAINTEXT://0.0.0.0:9092 等 → 客户端监听端口
  • KAFKA_CONTROLLER_QUORUM_VOTERS → 声明所有 Controller 节点,格式 ID@host:port

之前的内容存在错误,还是不能相信ai

修改要点

  • KAFKA_PROCESS_ROLES: 'broker,controller': 使用 Docker 服务名(kafka1、kafka2、kafka3),这样容器间通信时,Follower 才能正确找到 Leader。如果写成 localhost:9092,容器内会指向自己,副本同步失败。

  • 容器内端口的监听
    kafka1 监听 9092(与宿主机映射一致)
    kafka2 监听 9094(宿主机映射 9094:9094)
    kafka3 监听 9095(宿主机映射 9095:9095)

  • 3.Controller通信:使用9o93端口进行通信, 通过 KAFKA_CONTROLLER_QUORUM_VOTERS 中的服务名相互通信。

    4.3 advertised.listeners 说明

    在 Docker Compose 搭建的 Kafka 集群中,KAFKA_ADVERTISED_LISTENERS 的取值直接决定了 Broker 间通信的寻址方式。常见有两种写法:

    配置方式示例通信路径
    使用 Docker 服务名 PLAINTEXT://kafka1:9092 容器直连,通过 Compose 内部 DNS 解析
    使用宿主机物理 IP PLAINTEXT://192.168.100.168:9092 容器 → 宿主机 IP → 端口映射 → 目标容器

    两种方式都能工作吗?

    在大多数默认的 Docker Compose 网络环境下,两种配置都可以让集群启动并完成副本同步,因此你可能会在日志中看到 ISR 正常、Topic 收发无误。但两者的可靠性、性能和可移植性存在显著差异。

    服务名的优点

  • 通信路径更短
    容器间直接通过 Docker 内网通信,不经过宿主机的网络栈,延迟更低,吞吐量更高。

  • 不依赖宿主机 IP
    宿主机 IP 一旦变化(如 DHCP 重新分配、更换网络环境),集群将立即瘫痪,必须修改配置并重启所有节点。而服务名由 Compose 自动维护,永远指向正确的容器。

  • 可移植性强
    同一份 docker-compose.yml 可以直接在任何 Docker 主机上运行,无需修改 IP 地址。这对于分享教程、团队协作或 CI/CD 流水线至关重要。

  • 避免防火墙干扰
    容器直连不经过宿主机防火墙,而通过宿主机 IP 访问时可能被 iptables 或 firewalld 规则拦截,增加排查成本。

  • 什么时候可以用宿主机 IP?

    • 你明确需要从外部网络(非宿主机、非容器内)访问 Kafka,但又不想增加额外监听器。
    • 你只是在本地做一次性测试,且能接受上述局限性。
    如果既要内部高效通信,又要外部访问怎么办?

    Kafka 支持配置多个监听器,例如:

    • PLAINTEXT_INTERNAL://kafka1:9092 用于容器间复制
    • PLAINTEXT_EXTERNAL://192.168.100.168:9094 用于宿主机或外部客户端连接

    这个做法较复杂。在学习阶段,直接使用宿主机IP是更为简单,外部链接更高效的做法。


    五、启动集群

    5.1 创建数据目录并设置权限

    容器内 Kafka 进程以 appuser(UID=1000)运行,必须确保挂载目录可写。

    mkdir -p ./data/kafka1 ./data/kafka2 ./data/kafka3
    chown -R 1000:1000 ./data

    5.2 启动所有服务

    docker compose up -d # V2
    # 或 docker-compose up -d (V1)

    5.3 验证集群

    # 1. 容器状态
    docker compose ps

    # 2. 查看 kafka1 启动日志(搜索关键信息)
    docker logs kafka1 2>&1 | grep "Kafka Server started"
    # 如果没找到,直接看最后20行
    docker logs kafka1 –tail 20

    # 3. 查看 KRaft 元数据仲裁状态
    docker exec -it kafka1 kafka-metadata-quorum –bootstrap-server localhost:9092 describe

    正常输出示例:

    LeaderId: 1
    Voters: [1,2,3]
    Observers: []

    5.4 创建 Topic 并收发消息

    # 创建 Topic(3分区,3副本)
    docker exec -it kafka1 kafka-topics –create \\
    –topic test-topic \\
    –bootstrap-server localhost:9092 \\
    –partitions 3 \\
    –replication-factor 3

    # 查看 Topic 分区分布
    docker exec -it kafka1 kafka-topics –describe –topic test-topic –bootstrap-server localhost:9092

    # 生产消息
    docker exec -it kafka1 kafka-console-producer –topic test-topic –bootstrap-server localhost:9092,localhost:9094,localhost:9095

    # 消费消息(从头开始)
    docker exec -it kafka1 kafka-console-consumer –topic test-topic –from-beginning –bootstrap-server localhost:9092,localhost:9094,localhost:9095


    5.5 验证副本同步

    副本是否同步才是集群是否可靠的核心

    # 查看 Topic 详情
    docker exec -it kafka1 kafka-topics –describe –topic test –bootstrap-server kafka1:9092

    输出示例:

    Topic: testPartitionCount: 3ReplicationFactor: 3
    Partition: 0Leader: 1Replicas: 1,2,3Isr: 1,2,3
    Partition: 1Leader: 2Replicas: 2,3,1Isr: 2,3,1
    Partition: 2Leader: 3Replicas: 3,1,2Isr: 3,1,2

    如果 ISR 列始终为 1,2,3 全体成员,说明副本同步正常。

    如果 ISR 收缩到只剩 Leader(例如 Isr: 1),说明 KAFKA_ADVERTISED_LISTENERS 配置有误,请回头检查是否使用了服务名、端口是否匹配。

    六、常用命令

    以下命令还是很实用的,可以进行收藏。

    6.1 容器与日志管理

    # 查看容器运行状态
    docker compose ps

    # 查看所有容器实时资源占用
    docker stats

    # 查看 kafka1 最新 50 行日志
    docker logs kafka1 –tail 50

    # 实时跟踪 kafka2 日志
    docker logs -f kafka2

    # 搜索日志中是否包含特定内容(例如 ERROR)
    docker logs kafka1 2>&1 | grep ERROR

    # 查看最近 5 分钟内产生的日志(需使用 –since)
    docker logs –since 5m kafka3 | tail -100

    # 导出日志到文件进行分析
    docker logs kafka1 > kafka1.log 2>&1

    6.2 集群与元数据检查

    # 查询 KRaft 仲裁状态
    docker exec -it kafka1 kafka-metadata-quorum –bootstrap-server localhost:9092 describe

    # 查看所有 Broker 的版本和 API 信息
    docker exec -it kafka1 kafka-broker-api-versions –bootstrap-server localhost:9092

    # 列出所有 Topic
    docker exec -it kafka1 kafka-topics –list –bootstrap-server localhost:9092

    # 查看某个 Topic 详细信息(分区、副本分布、ISR)
    docker exec -it kafka1 kafka-topics –describe –topic test-topic –bootstrap-server localhost:9092

    # 查看消费者组列表
    docker exec -it kafka1 kafka-consumer-groups –bootstrap-server localhost:9092 –list

    # 查看特定消费者组的消费详情(lag 等)
    docker exec -it kafka1 kafka-consumer-groups –bootstrap-server localhost:9092 \\
    –group my-group –describe

    6.3 Topic 与分区管理

    # 创建 Topic(带压缩)
    docker exec -it kafka1 kafka-topics –create –topic compact-topic \\
    –bootstrap-server localhost:9092 –partitions 6 –replication-factor 2 \\
    –config cleanup.policy=compact

    # 修改分区数(只能增加)
    docker exec -it kafka1 kafka-topics –alter –topic test-topic \\
    –bootstrap-server localhost:9092 –partitions 6

    # 删除 Topic
    docker exec -it kafka1 kafka-topics –delete –topic old-topic \\
    –bootstrap-server localhost:9092

    # 查看 Topic 的消息数(需要消费工具)
    # 获得起始和结束偏移量
    docker exec -it kafka1 kafka-run-class kafka.tools.GetOffsetShell \\
    –broker-list localhost:9092 –topic test-topic –time -1

    6.4 生产者/消费者性能测试

    # 生产者压力测试(发送 10000 条记录,吞吐量 10000 msg/sec)
    docker exec -it kafka1 kafka-producer-perf-test \\
    –topic test-topic –num-records 10000 –record-size 100 \\
    –throughput 10000 –producer-props bootstrap.servers=localhost:9092

    # 消费者性能测试
    docker exec -it kafka1 kafka-consumer-perf-test \\
    –broker-list localhost:9092 –topic test-topic –messages 10000

    6.5 配置与动态管理

    # 查看 Broker 动态配置
    docker exec -it kafka1 kafka-configs –bootstrap-server localhost:9092 \\
    –broker 1 –describe –all

    # 动态设置 Broker 日志保留时间(1小时)
    docker exec -it kafka1 kafka-configs –bootstrap-server localhost:9092 \\
    –broker 1 –alter –add-config log.retention.hours=1

    # 查看 Topic 配置
    docker exec -it kafka1 kafka-configs –bootstrap-server localhost:9092 \\
    –topic test-topic –describe –all

    6.6 集群启停与清理

    # 优雅停止集群
    docker compose stop

    # 再次启动
    docker compose start

    # 彻底销毁所有容器和网络(保留数据目录)
    docker compose down

    # 销毁并删除所有数据卷(危险,数据将丢失)
    docker compose down -v

    6.7 其他实用命令

    # 进入容器内部操作
    docker exec -it kafka1 bash

    # 查看 JVM 堆内存使用(需要 jcmd,需进入容器)
    docker exec -it kafka1 jcmd 1 GC.heap_info

    # 查看 Kafka 进程 PID(容器内)
    docker exec -it kafka1 jps

    # 快速检查 Broker 是否在监听
    netstat -tlnp | grep 9092


    七、常见故障排查手册

    问题 1:端口冲突 port is already allocated

    # 定位占用进程
    sudo netstat -tlnp | grep -E '9092|9094|9095'
    # 若为旧容器,停止并删除
    docker stop <容器名> && docker rm <容器名>
    # 或直接清理整个旧 Compose 项目(在其目录下执行 docker compose down)

    问题 2:容器反复重启(日志显示 Restarting)

    首先查日志:docker logs kafka1 –tail 100

    日志关键词原因解决方法
    … writable FAILED 目录权限错误 chown -R 1000:1000 ./data
    does not appear to be a valid UUID CLUSTER_ID 非法 使用 Base64 UUID,并彻底清理数据:docker compose down rm -rf ./data mkdir -p ./data/kafka1 ./data/kafka2 ./data/kafka3 chown -R 1000:1000 ./data docker compose up -d
    Connection refused 或 Failed to connect to kafka1:9093 节点间网络不通或 Controller 监听器配置错误 检查 LISTENERS 是否包含 CONTROLLER://0.0.0.0:9093;确认防火墙未拦截容器内部网络

    问题 3:kafka-metadata-quorum 报错 unrecognized arguments: '–describe'

    正确写法:kafka-metadata-quorum –bootstrap-server localhost:9092 describe

    问题 4:Topic 创建失败,提示 No broker available

    • 等待 1~2 分钟,确保所有节点完全启动
    • 检查 KAFKA_ADVERTISED_LISTENERS 是否可被客户端解析
    • 如果 IP 发生变化,清理数据目录重新初始化

    问题 5:内存溢出(OOM)

    # 限制 Kafka 堆内存(在 docker-compose.yml 中添加)
    environment:
    KAFKA_HEAP_OPTS: "-Xmx512m -Xms512m"


    八、总结

    CLUSTER_ID 必须用 Base64 UUID → 用 kafka-storage random-uuid 生成

    KAFKA_ADVERTISED_LISTENERS 必须用服务名 → 容器间通信依赖 DNS 解析,localhost 会导致副本同步失败

    容器内监听端口要与宿主机映射端口对齐 → 否则宿主机客户端无法连接

    挂载目录权限设置为 1000:1000 → 容器内用户需要写权限

    按照本文的最终版配置,你可以搭建出一个真正健康、副本同步正常的 KRaft Kafka 集群。后续可以在此基础上添加监控(如 Prometheus)、修改为多机部署,或集成到 Spring Boot 等微服务中。

    如果本文对你有帮助,欢迎点赞、收藏、关注,后续会出一个借助AI学习kafka相关入门知识的笔记
    欢迎大家批评指正,也欢迎在评论区交流你遇到的问题!


    最后更新:2026-05-30

    赞(0)
    未经允许不得转载:171主机测评 » Kafka入门集群搭建-学习指南
    分享到: 更多 (0)

    评论 抢沙发

    • 昵称 (必填)
    • 邮箱 (必填)
    • 网址