Docker环境下搭建Kafka集群完整实战教程(含生产者消费者测试)
📝 前言
Kafka作为目前最流行的分布式消息队列系统,在大数据实时处理场景中扮演着至关重要的角色。本文将手把手带你在Docker环境下搭建一个3节点的Kafka集群,并通过实际案例验证集群的消息传输能力。
适合人群:
- 初学Kafka的开发者
- 学习实时计算的学生
技术栈:
- Docker & Docker Compose
- Apache Kafka 3.6.1
- Kraft
🎯 实验目标
- 理解Kafka的核心组件和工作原理
- 掌握在Docker容器中部署Kafka集群的方法
- 验证Kafka集群的消息生产和消费功能
- 测试集群处理实时数据的性能
🚀 第一步:拉取Kafka镜像并创建网络
1.1 拉取Bitnami Kafka镜像
首先启动Docker,然后在命令行中执行:
docker pull bitnami/kafka:3.6.1
等待镜像下载完成,你会看到类似下面的成功提示。
1.2 创建自定义Docker网络
为了让Kafka集群各节点之间能够相互通信,我们需要创建一个自定义网络:
docker network create --subnet=172.23.0.0/16 --gateway=172.23.0.1 netkafka
参数说明:
--subnet: 指定网络的IP地址范围--gateway: 指定网关地址netkafka: 自定义网络名称
⚙️ 第二步:配置并启动Kafka集群
2.1 准备docker-compose配置文件
创建一个名为 kafka.yml 的文件,配置3个Kafka节点。以下是完整配置:
version: "3.6"
services:
kafka1:
container_name: kafka1
image: 'bitnami/kafka:3.6.1'
user: root
ports:
- '19092:9092'
- '19093:9093'
environment:
# 允许使用Kraft
- KAFKA_ENABLE_KRAFT=yes
- KAFKA_CFG_PROCESS_ROLES=broker,controller
- KAFKA_CFG_CONTROLLER_LISTENER_NAMES=CONTROLLER
# 定义kafka服务端socket监听端口(Docker内部的ip地址和端口)
- KAFKA_CFG_LISTENERS=PLAINTEXT://:9092,CONTROLLER://:9093
# 定义安全协议
- KAFKA_CFG_LISTENER_SECURITY_PROTOCOL_MAP=CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT
#定义外网访问地址(宿主机ip地址和端口,标红处修改为自己主机IP)
- KAFKA_CFG_ADVERTISED_LISTENERS=PLAINTEXT://10.86.81.70:19092
- KAFKA_CFG_NODE_ID=1
- KAFKA_KRAFT_CLUSTER_ID=iZWRiSqjZAlYwlKEqHFQWI
- KAFKA_CFG_CONTROLLER_QUORUM_VOTERS=1@172.23.0.11:9093,2@172.23.0.12:9093,3@172.23.0.13:9093
- ALLOW_PLAINTEXT_LISTENER=yes
# 设置broker最大内存,和初始内存
- KAFKA_HEAP_OPTS=-Xmx512M -Xms256M
volumes:
#挂载路径,标红处修改为自己的路径
- D:/Realtimecomputation/dockers/volume/kafka/broker01:/bitnami/kafka:rw
networks:
netkafka:
ipv4_address: 172.23.0.11
kafka2:
container_name: kafka2
image: 'bitnami/kafka:3.6.1'
user: root
ports:
- '29092:9092'
- '29093:9093'
environment:
- KAFKA_ENABLE_KRAFT=yes
- KAFKA_CFG_PROCESS_ROLES=broker,controller
- KAFKA_CFG_CONTROLLER_LISTENER_NAMES=CONTROLLER
- KAFKA_CFG_LISTENERS=PLAINTEXT://:9092,CONTROLLER://:9093
- KAFKA_CFG_LISTENER_SECURITY_PROTOCOL_MAP=CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT
#标红处修改为自己主机IP
- KAFKA_CFG_ADVERTISED_LISTENERS=PLAINTEXT://10.86.81.70:29092
- KAFKA_CFG_NODE_ID=2
- KAFKA_KRAFT_CLUSTER_ID=iZWRiSqjZAlYwlKEqHFQWI #哪一,三个节点保持一致
- KAFKA_CFG_CONTROLLER_QUORUM_VOTERS=1@172.23.0.11:9093,2@172.23.0.12:9093,3@172.23.0.13:9093
- ALLOW_PLAINTEXT_LISTENER=yes
- KAFKA_HEAP_OPTS=-Xmx512M -Xms256M
volumes:
#挂载路径,标红处修改为自己的路径
- D:/Realtimecomputation/dockers/volume/kafka/broker02:/bitnami/kafka:rw
networks:
netkafka:
ipv4_address: 172.23.0.12
kafka3:
container_name: kafka3
image: 'bitnami/kafka:3.6.1'
user: root
ports:
- '39092:9092'
- '39093:9093'
environment:
- KAFKA_ENABLE_KRAFT=yes
- KAFKA_CFG_PROCESS_ROLES=broker,controller
- KAFKA_CFG_CONTROLLER_LISTENER_NAMES=CONTROLLER
- KAFKA_CFG_LISTENERS=PLAINTEXT://:9092,CONTROLLER://:9093
- KAFKA_CFG_LISTENER_SECURITY_PROTOCOL_MAP=CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT
# 标红处修改为自己主机IP
- KAFKA_CFG_ADVERTISED_LISTENERS=PLAINTEXT://10.86.81.70:39092
- KAFKA_CFG_NODE_ID=3
- KAFKA_KRAFT_CLUSTER_ID=iZWRiSqjZAlYwlKEqHFQWI
- KAFKA_CFG_CONTROLLER_QUORUM_VOTERS=1@172.23.0.11:9093,2@172.23.0.12:9093,3@172.23.0.13:9093
- ALLOW_PLAINTEXT_LISTENER=yes
- KAFKA_HEAP_OPTS=-Xmx512M -Xms256M
volumes:
#挂载路径,标红处修改为自己的路径
- D:/Realtimecomputation/dockers/volume/kafka/broker03:/bitnami/kafka:rw
networks:
netkafka:
ipv4_address: 172.23.0.13
networks:
name:
netkafka:
external: true
driver: bridge
name: netkafka
ipam:
driver: default
config:
- subnet: 172.23.0.0/25
gateway: 172.23.0.1
2.2 重要配置说明
必须修改的两个地方:
-
宿主机IP地址 -
KAFKA_CFG_ADVERTISED_LISTENERS- 将
172.20.10.5改为你的本机IP地址 - 如何查看本机IP:在cmd中输入
ipconfig
- 将
-
挂载路径 -
volumes- 将
D:/Realtimecomputation/...改为你的本地路径 - 确保该路径已存在且有读写权限
- 将
2.3 启动Kafka集群
在 kafka.yml 文件所在目录执行:
docker-compose -f kafka.yml up -d

启动成功后,打开Docker Desktop,你会看到3个Kafka容器正在运行。

🔬 第三步:测试生产者和消费者
3.1 创建Topic
进入kafka1容器:
docker exec -it kafka1 bash
cd /opt/bitnami/kafka/bin
创建一个名为 foo 的Topic,设置5个分区和3个副本:
./kafka-topics.sh --create --topic foo --partitions 5 \
--replication-factor 3 \
--bootstrap-server kafka1:9092,kafka2:9092,kafka3:9092
查看Topic详细信息:
./kafka-topics.sh --describe --topic foo \
--bootstrap-server kafka1:9092,kafka2:9092,kafka3:9092
你会看到分区分布、副本信息等详细数据。
3.2 创建生产者
打开一个新的cmd窗口,进入kafka1容器:
docker exec -it kafka1 bash
创建生产者:
kafka-console-producer.sh --broker-list \
172.23.0.11:9092,172.23.0.12:9092,172.23.0.13:9092 \
--topic foo
现在可以在这个窗口输入消息了!
3.3 创建消费者(两个)
消费者1 - 在kafka2容器中:
# 打开新cmd窗口
docker exec -it kafka2 bash
# 创建消费者
kafka-console-consumer.sh --bootstrap-server \
172.23.0.11:9092,172.23.0.12:9092,172.23.0.13:9092 \
--topic foo
消费者2 - 在kafka3容器中:
# 再打开一个新cmd窗口
docker exec -it kafka3 bash
# 创建消费者
kafka-console-consumer.sh --bootstrap-server \
172.23.0.11:9092,172.23.0.12:9092,172.23.0.13:9093 \
--topic foo
3.4 验证消息传输
在生产者窗口输入消息:
> hello
> 111
> test message
你会看到两个消费者窗口都能实时接收到这些消息!这证明Kafka集群的消息传输功能正常。
📊 第四步:接入实时数据测试
4.1 准备数据文件
将CSV数据文件复制到kafka1容器:
docker cp D:\Realtimecomputation\lab2-data\stock-part1.csv kafka1:/tmp/
docker cp D:\Realtimecomputation\lab2-data\stock-part2.csv kafka1:/tmp/
4.2 批量发送数据
在kafka1容器中执行:
cat /tmp/stock-part1.csv | kafka-console-producer.sh \
--broker-list 172.23.0.11:9092,172.23.0.12:9092,172.23.0.13:9092 \
--topic foo
4.3 观察消费情况
此时两个消费者会实时接收到CSV文件中的所有数据,证明Kafka集群能够高效处理批量实时数据。

📈 实验结果分析
核心成果
- 集群稳定性:3节点Kafka集群成功部署并稳定运行,实现了高可用架构
- 消息传输:生产者发送的消息能被不同节点的消费者及时接收,无丢失
- 实时处理:批量导入CSV数据时,两个消费者均能正常接收,无明显延迟
- 副本机制:通过3副本配置,保证了数据的可靠性和容错能力
性能特点
- 低延迟:消息从生产到消费几乎是实时的
- 高吞吐:能够快速处理批量数据
- 容错性:即使单个节点故障,集群仍可正常工作(未在本次实验中测试)
⚠️ 常见问题及解决方案
问题1:集群连接失败
原因: 宿主机IP地址配置错误
解决方案:
- 使用
ipconfig查看当前IP地址 - 修改
kafka.yml中所有的KAFKA_CFG_ADVERTISED_LISTENERS - 重启集群:
docker-compose -f kafka.yml restart
问题2:消费者接收到乱码
原因: CSV文件编码格式不是UTF-8
解决方案:
- 用Excel打开CSV文件
- 另存为时选择UTF-8编码
- 重新导入数据
问题3:容器启动失败
可能原因:
- 端口被占用
- 挂载路径不存在
- Docker资源不足
解决方案:
- 检查端口占用:
netstat -ano | findstr "19092" - 确保挂载路径存在且有权限
- 分配更多内存给Docker
🎓 总结
通过本次实战,我们完成了:
✅ 使用Docker快速搭建Kafka集群
✅ 理解Kafka的核心概念(Topic、Partition、Replica)
✅ 掌握生产者和消费者的基本使用
✅ 验证了集群的实时数据处理能力
Kafka作为大数据生态系统的重要组件,是学习实时计算必须掌握的技术。这个实验为后续学习Flink、Spark等框架与Kafka的协同使用打下了坚实基础。
更多推荐
所有评论(0)