Kafka 3.1.0单机与集群环境搭建指南

1. Kafka环境搭建概述

第一次接触Kafka时,我被它"分布式消息队列"的名头吓到了,直到真正动手搭建才发现其实并不复杂。Kafka作为当前最流行的分布式消息系统,在日志收集、流处理、事件溯源等场景都有广泛应用。本文将基于Kafka 3.1.0版本,带你从零开始完成单机和集群两种环境的搭建。

注意:虽然Kafka官方已开始推荐使用KRaft模式(去ZooKeeper化),但考虑到生产环境多数仍在使用ZooKeeper协调模式,本文仍采用传统架构进行演示。

2. 基础环境准备

2.1 JDK安装与配置

Kafka基于Scala开发,运行需要Java环境。根据官方文档,Kafka 3.x支持Java 8和11,这里以Java 8为例:

# 下载JDK (以jdk-8u291为例) wget https://download.oracle.com/java/8u291/b09/jdk-8u291-linux-x64.tar.gz # 解压到/usr/local目录 tar -zxvf jdk-8u291-linux-x64.tar.gz -C /usr/local/

配置环境变量时有个小技巧:在/etc/profile.d/目录下单独创建java.sh,这样系统更新时不会被覆盖:

# /etc/profile.d/java.sh export JAVA_HOME=/usr/local/jdk1.8.0_291 export PATH=$JAVA_HOME/bin:$PATH export CLASSPATH=.:$JAVA_HOME/lib/dt.jar:$JAVA_HOME/lib/tools.jar

执行source /etc/profile使配置生效后,用java -version验证安装:

java version "1.8.0_291" Java(TM) SE Runtime Environment (build 1.8.0_291-b09)

2.2 Kafka安装包获取

建议直接从Apache官网下载二进制包,避免源码编译的兼容性问题:

wget https://archive.apache.org/dist/kafka/3.1.0/kafka_2.13-3.1.0.tgz tar -xvf kafka_2.13-3.1.0.tgz -C /usr/local/ cd /usr/local/kafka_2.13-3.1.0

实测发现:2.13表示Scala版本,3.1.0是Kafka版本,版本不匹配会导致运行时出现奇怪的ClassNotFound错误。

3. 单机模式部署

3.1 ZooKeeper服务启动

当前版本Kafka仍依赖ZooKeeper做元数据管理,先启动内置的ZooKeeper:

# 后台启动ZooKeeper nohup bin/zookeeper-server-start.sh config/zookeeper.properties > zookeeper.log 2>&1 &

检查是否启动成功:

netstat -tunlp | grep 2181 tcp6 0 0 :::2181 :::* LISTEN 12345/java

3.2 Kafka服务配置

修改config/server.properties关键参数:

# 必须设置真实IP,否则远程客户端无法连接 listeners=PLAINTEXT://192.168.1.100:9092 # 日志保存时间(小时) log.retention.hours=168 # 单个日志文件大小 log.segment.bytes=1073741824

启动Kafka服务:

nohup bin/kafka-server-start.sh config/server.properties > kafka.log 2>&1 &

验证服务状态:

jps 12345 QuorumPeerMain # ZooKeeper进程 23456 Kafka # Kafka进程

4. 基础功能测试

4.1 Topic管理

创建测试Topic(1个分区,2个副本):

bin/kafka-topics.sh --create \ --topic test-topic \ --bootstrap-server 192.168.1.100:9092 \ --partitions 1 \ --replication-factor 1

查看Topic详情:

bin/kafka-topics.sh --describe \ --topic test-topic \ --bootstrap-server 192.168.1.100:9092

输出示例:

Topic: test-topic PartitionCount: 1 ReplicationFactor: 1 Configs: Topic: test-topic Partition: 0 Leader: 0 Replicas: 0 Isr: 0

4.2 生产者消费者测试

开两个终端分别运行:

# 生产者 bin/kafka-console-producer.sh \ --topic test-topic \ --bootstrap-server 192.168.1.100:9092 # 消费者(从头开始消费) bin/kafka-console-consumer.sh \ --topic test-topic \ --from-beginning \ --bootstrap-server 192.168.1.100:9092

在生产者终端输入消息,消费者终端应能实时接收到。

5. 集群模式部署

5.1 ZooKeeper集群配置

准备3台服务器(192.168.1.101-103),每台修改config/zookeeper.properties:

dataDir=/var/lib/zookeeper clientPort=2181 initLimit=5 syncLimit=2 server.1=192.168.1.101:2888:3888 server.2=192.168.1.102:2888:3888 server.3=192.168.1.103:2888:3888

在各节点创建myid文件:

# 节点1 echo 1 > /var/lib/zookeeper/myid # 节点2 echo 2 > /var/lib/zookeeper/myid # 节点3 echo 3 > /var/lib/zookeeper/myid

启动ZooKeeper集群:

bin/zookeeper-server-start.sh config/zookeeper.properties

5.2 Kafka集群配置

各节点修改config/server.properties:

# 节点1配置 broker.id=1 listeners=PLAINTEXT://192.168.1.101:9092 zookeeper.connect=192.168.1.101:2181,192.168.1.102:2181,192.168.1.103:2181 # 节点2配置 broker.id=2 listeners=PLAINTEXT://192.168.1.102:9092 # 节点3配置 broker.id=3 listeners=PLAINTEXT://192.168.1.103:9092

启动所有Kafka节点:

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

5.3 集群验证

创建Topic(3分区,2副本):

bin/kafka-topics.sh --create \ --topic cluster-test \ --bootstrap-server 192.168.1.101:9092 \ --partitions 3 \ --replication-factor 2

查看分区分布情况:

bin/kafka-topics.sh --describe \ --topic cluster-test \ --bootstrap-server 192.168.1.101:9092

正常输出应类似:

Topic: cluster-test PartitionCount: 3 ReplicationFactor: 2 Configs: Topic: cluster-test Partition: 0 Leader: 2 Replicas: 2,1 Isr: 2,1 Topic: cluster-test Partition: 1 Leader: 3 Replicas: 3,2 Isr: 3,2 Topic: cluster-test Partition: 2 Leader: 1 Replicas: 1,3 Isr: 1,3

6. 常见问题排查

6.1 连接问题

错误现象:

Connection to node -1 could not be established. Broker may not be available.

解决方案:

  1. 检查server.properties中的listeners配置
  2. 确认防火墙开放了9092端口
  3. 测试telnet IP 9092验证网络连通性

6.2 ZooKeeper连接问题

错误日志:

ERROR [KafkaServer id=1] Fatal error during KafkaServer startup (kafka.server.KafkaServer) org.apache.zookeeper.KeeperException$ConnectionLossException: KeeperErrorCode = ConnectionLoss

解决方法:

  1. 检查ZooKeeper服务状态
  2. 确认zookeeper.connect配置正确
  3. 查看ZooKeeper日志排查具体原因

6.3 磁盘空间不足

错误日志:

ERROR [Log partition=test-0, dir=/tmp/kafka-logs] Error while loading log dir /tmp/kafka-logs/test-0 (kafka.log.Log) java.io.IOException: No space left on device

建议方案:

  1. 修改config/server.properties中的log.dirs参数
  2. 设置自动清理策略:
    log.retention.hours=168 log.retention.bytes=10737418240

7. 生产环境建议

  1. 日志目录分离:将Kafka日志与系统盘分离,避免影响系统运行

    log.dirs=/data/kafka-logs
  2. JVM调优:根据服务器配置调整内存

    export KAFKA_HEAP_OPTS="-Xms4G -Xmx4G"
  3. 监控配置:建议启用JMX监控

    JMX_PORT=9999
  4. 安全加固:生产环境应配置SASL认证

    security.inter.broker.protocol=SASL_PLAINTEXT sasl.mechanism.inter.broker.protocol=PLAIN

我在实际部署中发现,Kafka对文件描述符数要求较高,建议提前调整:

ulimit -n 100000 echo "* soft nofile 100000" >> /etc/security/limits.conf