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

跨两台VPS的Apache Kafka单向通信异常排查求助

双向Kafka通信故障排查与解决

问题背景

我在两台VPS(IP分别为x.x.x.x和y.y.y.y)上测试Apache Kafka,两台机器均部署了ZooKeeper和Kafka Server,生产者、消费者代码逻辑一致,具体代码如下:

第一台VPS(x.x.x.x)

Producer.py

from kafka import KafkaProducer
import json
# Initialize a Kafka Producer
producer = KafkaProducer(bootstrap_servers=['x.x.x.x:9092'],
                         api_version=(0,10),
                         value_serializer=lambda x: json.dumps(x).encode('utf-8'))

# Send a message
message = {"data": "from first message"}
producer.send('first', value=message)

# Ensure all messages are sent
producer.flush()

Consumer.py

from kafka import KafkaConsumer
import json

# Initialize a Kafka Consumer
consumer = KafkaConsumer('second',
                         bootstrap_servers=['y.y.y.y:9092'],
                         auto_offset_reset='earliest',
                         api_version=(0, 10),
                         value_deserializer=lambda x: json.loads(x.decode('utf-8')))

# Print received messages
for message in consumer:
    print(message.value)

第二台VPS(y.y.y.y)

Producer.py

from kafka import KafkaProducer
import json
# Initialize a Kafka Producer
producer = KafkaProducer(bootstrap_servers=['y.y.y.y:9092'],
                         api_version=(0,10),
                         value_serializer=lambda x: json.dumps(x).encode('utf-8'))

# Send a message
message = {"data": "from second message"}
producer.send('second', value=message)

# Ensure all messages are sent
producer.flush()

Consumer.py

from kafka import KafkaConsumer
import json

# Initialize a Kafka Consumer
consumer = KafkaConsumer('first',
                         bootstrap_servers=['x.x.x.x:9092'],
                         auto_offset_reset='earliest',
                         api_version=(0, 10),
                         value_deserializer=lambda x: json.loads(x.decode('utf-8')))

# Print received messages
for message in consumer:
    print(message.value)

当前状态

  • 运行x.x.x.x上的生产者,y.y.y.y上的消费者可正常接收消息
  • 运行y.y.y.y上的生产者,x.x.x.x上的消费者无法接收任何消息

由于代码逻辑完全一致,判断问题出在Apache Kafka Server配置上,需调整配置实现双向通信。

解决方案

1. 修正advertised.listeners与listeners配置

Kafka的advertised.listeners是客户端实际使用的连接地址,必须设置为VPS的公网IP,不能用localhost或内网IP。两台机器的server.properties都需修改:

第一台VPS(x.x.x.x)配置

advertised.listeners=PLAINTEXT://x.x.x.x:9092
listeners=PLAINTEXT://0.0.0.0:9092

第二台VPS(y.y.y.y)配置

advertised.listeners=PLAINTEXT://y.y.y.y:9092
listeners=PLAINTEXT://0.0.0.0:9092
  • listeners=0.0.0.0:9092让Kafka监听所有网卡,确保能接收外部连接
  • advertised.listeners必须配置为客户端可直接访问的IP,否则生产者/消费者会获取错误的连接地址

2. 检查并开放端口

确保两台VPS的9092端口(Kafka默认端口)对外开放,防火墙允许双向访问:

  • Ubuntu系统:
    # 查看防火墙状态
    sudo ufw status
    # 开放9092端口
    sudo ufw allow 9092
    
  • CentOS/RHEL系统:
    # 查看防火墙规则
    firewall-cmd --list-all
    # 开放端口并重启防火墙
    firewall-cmd --add-port=9092/tcp --permanent
    firewall-cmd --reload
    

3. 验证ZooKeeper集群配置(若为集群模式)

如果两台ZooKeeper组成集群,确保zoo.cfg中正确配置彼此地址:

server.1=x.x.x.x:2888:3888
server.2=y.y.y.y:2888:3888

同时在每台机器的ZooKeeper数据目录(dataDir指定路径)下创建myid文件,第一台写入1,第二台写入2,然后重启ZooKeeper服务。

4. 重启服务并验证

修改配置后,依次重启ZooKeeper和Kafka服务:

# 停止服务
sudo systemctl stop zookeeper
sudo systemctl stop kafka

# 启动服务
sudo systemctl start zookeeper
sudo systemctl start kafka

重启完成后重新测试生产者与消费者的双向通信。

额外排查步骤

  • 用telnet测试端口连通性:在x.x.x.x执行telnet y.y.y.y 9092,在y.y.y.y执行telnet x.x.x.x 9092,确保端口能正常连通
  • 查看Kafka日志(默认路径/var/log/kafka/),检查是否存在连接错误或配置异常的日志信息

内容的提问来源于stack exchange,提问作者Andy Wood

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.26 12:15:56