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/java3.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: 04.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.properties5.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.properties5.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,36. 常见问题排查
6.1 连接问题
错误现象:
Connection to node -1 could not be established. Broker may not be available.解决方案:
- 检查server.properties中的listeners配置
- 确认防火墙开放了9092端口
- 测试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解决方法:
- 检查ZooKeeper服务状态
- 确认zookeeper.connect配置正确
- 查看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建议方案:
- 修改config/server.properties中的log.dirs参数
- 设置自动清理策略:
log.retention.hours=168 log.retention.bytes=10737418240
7. 生产环境建议
日志目录分离:将Kafka日志与系统盘分离,避免影响系统运行
log.dirs=/data/kafka-logsJVM调优:根据服务器配置调整内存
export KAFKA_HEAP_OPTS="-Xms4G -Xmx4G"监控配置:建议启用JMX监控
JMX_PORT=9999安全加固:生产环境应配置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