📝 前言

        Kafka作为目前最流行的分布式消息队列系统,在大数据实时处理场景中扮演着至关重要的角色。本文将手把手带你在Docker环境下搭建一个3节点的Kafka集群,并通过实际案例验证集群的消息传输能力。

适合人群:

  • 初学Kafka的开发者
  • 学习实时计算的学生

技术栈:

  • Docker & Docker Compose
  • Apache Kafka 3.6.1
  • Kraft

🎯 实验目标

  1. 理解Kafka的核心组件和工作原理
  2. 掌握在Docker容器中部署Kafka集群的方法
  3. 验证Kafka集群的消息生产和消费功能
  4. 测试集群处理实时数据的性能

🚀 第一步:拉取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 重要配置说明

必须修改的两个地方:

  1. 宿主机IP地址 - KAFKA_CFG_ADVERTISED_LISTENERS

    • 172.20.10.5 改为你的本机IP地址
    • 如何查看本机IP:在cmd中输入 ipconfig
  2. 挂载路径 - 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集群能够高效处理批量实时数据。


📈 实验结果分析

核心成果

  1. 集群稳定性:3节点Kafka集群成功部署并稳定运行,实现了高可用架构
  2. 消息传输:生产者发送的消息能被不同节点的消费者及时接收,无丢失
  3. 实时处理:批量导入CSV数据时,两个消费者均能正常接收,无明显延迟
  4. 副本机制:通过3副本配置,保证了数据的可靠性和容错能力

性能特点

  • 低延迟:消息从生产到消费几乎是实时的
  • 高吞吐:能够快速处理批量数据
  • 容错性:即使单个节点故障,集群仍可正常工作(未在本次实验中测试)

⚠️ 常见问题及解决方案

问题1:集群连接失败

原因: 宿主机IP地址配置错误

解决方案:

  1. 使用 ipconfig 查看当前IP地址
  2. 修改 kafka.yml 中所有的 KAFKA_CFG_ADVERTISED_LISTENERS
  3. 重启集群:docker-compose -f kafka.yml restart

问题2:消费者接收到乱码

原因: CSV文件编码格式不是UTF-8

解决方案:

  1. 用Excel打开CSV文件
  2. 另存为时选择UTF-8编码
  3. 重新导入数据

问题3:容器启动失败

可能原因:

  • 端口被占用
  • 挂载路径不存在
  • Docker资源不足

解决方案:

  • 检查端口占用:netstat -ano | findstr "19092"
  • 确保挂载路径存在且有权限
  • 分配更多内存给Docker


🎓 总结

通过本次实战,我们完成了:

✅ 使用Docker快速搭建Kafka集群
✅ 理解Kafka的核心概念(Topic、Partition、Replica)
✅ 掌握生产者和消费者的基本使用
✅ 验证了集群的实时数据处理能力

Kafka作为大数据生态系统的重要组件,是学习实时计算必须掌握的技术。这个实验为后续学习Flink、Spark等框架与Kafka的协同使用打下了坚实基础。

更多推荐