You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

多节点Kafka集群连接失败且消息仅写入单分区问题求助

问题排查与解决:Kafka集群连接拒绝+分区集中异常

问题现象

  1. 连接拒绝错误:使用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)
  1. 消息分区写入异常:回调函数仅返回单分区写入结果:

Message delivered to topic : json_test [partition : 0],

  1. 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+对外端口,例如:
    advertised.listeners=PLAINTEXT://100.25.177.77:9095
    
    每个Broker需对应自己的宿主机IP和端口
  • listeners:配置为容器内部监听地址,例如:
    listeners=PLAINTEXT://0.0.0.0:9092
    
  • broker.id:确保每个Broker的ID唯一(如0、1、2)
  • 重启所有Broker后,用kafkacat验证集群状态:
    kafkacat -L -b <任意Broker地址:端口>
    
    正常情况下应能看到所有3个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地址:端口>
    
  • 方法二:使用分区重分配工具
    1. 创建重分配配置文件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]}
      ]
    }
    
    1. 执行重分配:
    kafka-reassign-partitions.sh --bootstrap-server <正常Broker地址:端口> --reassignment-json-file reassign.json --execute
    
    1. 验证重分配结果:
    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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.26 07:44:59