多节点Kafka集群连接失败且消息仅写入单分区问题求助
问题排查与解决:Kafka集群连接拒绝+分区集中异常
问题现象
- 连接拒绝错误:使用
confluent-kafka-python生产数据时,频繁出现以下错误:
> %3|1680100333.779|FAIL|rdkafka#producer-1| [thrd:100.25.177.77:9095/bootstrap]: 100.25.177.77:9095/bootstrap: Connect to ipv4#100.25.177.77:9095 failed: Connection refused (after 190ms in state CONNECT) > %3|1680100333.973|FAIL|rdkafka#producer-1| [thrd:54.152.58.40:9094/bootstrap]: 54.152.58.40:9094/bootstrap: Connect to ipv4#54.152.58.40:9094 failed: Connection refused (after 194ms in state CONNECT)
- 消息分区写入异常:回调函数仅返回单分区写入结果:
Message delivered to topic : json_test [partition : 0],
- Topic分区集中问题:创建的Topic所有分区均集中在单个Broker,创建Topic的代码如下:
new_topics = [NewTopic('ec2_json_test', num_partitions=3, replication_factor=1)]
多次执行脚本后,客户端连接的Broker随机,但Topic分区始终未分散到集群节点。
环境信息
- EC2实例:3台t2.micro
- Kafka版本:2.13-3.4.0
- 部署方式:Docker Swarm部署Kafka集群,ZooKeeper仅部署在管理节点;其他节点9094、9095端口均返回连接拒绝;
kafkacat仅能查询到1个Broker信息,所有Topic分区leader均为该Broker。
排查与解决步骤
1. 验证Broker端口可达性
- 在客户端机器上用
nc或telnet测试所有Broker的端口连通性,例如:nc -zv 100.25.177.77 9095 nc -zv 54.152.58.40 9094 - 检查EC2安全组:确保客户端IP(或0.0.0.0/0用于测试)能访问Broker宿主机的9094/9095等端口
- 检查Docker Swarm端口映射:确认Kafka容器的内部端口(如9092)已正确映射到宿主机的9094/9095端口,且容器处于运行状态
2. 修正Kafka Broker核心配置
Kafka集群无法被正常识别的核心原因通常是**advertised.listeners配置错误**,需针对每个Broker单独配置:
advertised.listeners:必须配置为Broker宿主机的公网/内网IP+对外端口,例如:
每个Broker需对应自己的宿主机IP和端口advertised.listeners=PLAINTEXT://100.25.177.77:9095listeners:配置为容器内部监听地址,例如:listeners=PLAINTEXT://0.0.0.0:9092broker.id:确保每个Broker的ID唯一(如0、1、2)- 重启所有Broker后,用
kafkacat验证集群状态:
正常情况下应能看到所有3个Broker信息kafkacat -L -b <任意Broker地址:端口>
3. 重新均衡Topic分区
当集群恢复正常后,修复现有Topic的分区分布:
- 方法一:删除并重建Topic
kafka-topics.sh --delete --topic ec2_json_test --bootstrap-server <正常Broker地址:端口> kafka-topics.sh --create --topic ec2_json_test --num-partitions 3 --replication-factor 1 --bootstrap-server <正常Broker地址:端口> - 方法二:使用分区重分配工具
- 创建重分配配置文件
reassign.json:
{ "topics": [{"topic": "ec2_json_test"}], "version": 1, "partitions": [ {"topic": "ec2_json_test", "partition": 0, "replicas": [0]}, {"topic": "ec2_json_test", "partition": 1, "replicas": [1]}, {"topic": "ec2_json_test", "partition": 2, "replicas": [2]} ] }- 执行重分配:
kafka-reassign-partitions.sh --bootstrap-server <正常Broker地址:端口> --reassignment-json-file reassign.json --execute- 验证重分配结果:
kafka-reassign-partitions.sh --bootstrap-server <正常Broker地址:端口> --reassignment-json-file reassign.json --verify - 创建重分配配置文件
4. 优化生产者配置
- 在Python生产者的配置中,
bootstrap.servers需配置所有Broker的地址,示例:producer_config = { 'bootstrap.servers': '100.25.177.77:9095,54.152.58.40:9094,<第三个BrokerIP>:<端口>', 'acks': 'all', 'key.serializer': lambda x: json.dumps(x).encode('utf-8'), 'value.serializer': lambda x: json.dumps(x).encode('utf-8') }
内容的提问来源于stack exchange,提问作者humanlearning
相关产品推荐
相关产品推荐

