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

远程运行Kafka Java Consumer遇连接拒绝问题求助

Kafka远程Consumer连接拒绝问题排查

问题描述

我把Kafka Server和Producer部署在系统A上,自己用Java写了一个Kafka Consumer:

  • 在系统A本地运行时,Consumer能正常接收消息;
  • 但放到其他系统(比如系统B)运行时,就报连接拒绝的异常,没法接收消息。
    已经确认系统B和A之间的ping、telnet都是通的。

关键日志显示:Consumer初始化时指定了bootstrap.servers=192.168.95.217:9092,也成功连接这个节点拿到了元数据,但之后却尝试去连localhost:9092,最后抛出java.net.ConnectException: Connection refused。

我的Java Consumer代码如下:

public class KafkaSubscription {
    public static void main(String args[]) throws InterruptedException {
        Properties props = new Properties();
        props.put("bootstrap.servers", "192.168.95.217:9092");
        // props.put("zookeeper.connect", "192.168.95.217:2181");
        props.put("group.id", "test-consumer-group1");
        props.put("enable.auto.commit", "true");
        props.put("auto.commit.interval.ms", "1000");
        props.put("auto.offset.reset", "earliest");
        props.put("session.timeout.ms", "30000");
        props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
        props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
        KafkaConsumer<String, String> kafkaConsumer = new KafkaConsumer<>(props);
        System.out.println("properties loaded");
        kafkaConsumer.subscribe(Arrays.asList("test1"));
        kafkaConsumer.seekToBeginning(Collections.emptyList());
        while (true) {
            ConsumerRecords<String, String> records = kafkaConsumer.poll(1);
            System.out.println("Records Length : " + records.count());
            for (ConsumerRecord<String, String> record : records) {
                System.out.printf("offset = %d, value = %s", record.offset(), record.value());
                System.out.println();
            }
            Thread.sleep(30000);
        }
    }
}

想请教下,还需要检查哪些配置或设置,才能让远程Consumer正常接收消息?


排查与解决建议

这种情况我碰到过好几次,核心原因是Kafka Broker对外暴露的地址配置不对——Consumer通过bootstrap server拿到元数据后,会按照Broker在元数据里返回的地址去建立连接,而不是一直用你配置的bootstrap.servers地址。你现在看到它去连localhost,说明Broker告诉Consumer自己的地址是localhost:9092,远程的系统B自然连不上。

你需要检查系统A上Kafka Broker的以下配置:

  • 检查advertised.listeners配置
    这是最关键的配置,决定了Broker对外暴露给客户端的地址。默认情况下,Kafka会用listeners的配置或者自动识别的地址,但如果是在多网卡或者远程访问场景下,必须手动设置这个参数。
    打开Kafka的配置文件(通常是server.properties),把advertised.listeners设置成系统A的局域网可访问的IP,比如:

    advertised.listeners=PLAINTEXT://192.168.95.217:9092
    

    注意要和你Consumer里配置的bootstrap.servers地址一致,或者是客户端能访问到的有效地址。

  • 确认listeners配置
    listeners是Broker监听的本地地址,一般建议设置成PLAINTEXT://0.0.0.0:9092,这样Broker会监听所有网卡的9092端口,确保能接收来自系统B的连接。如果这个配置只设了localhost:9092,那Broker只会监听本地请求,远程连过来也会失败——不过你说telnet能通,可能这个已经配置对了,但还是要确认下。

  • 检查listener.security.protocol.map(如果有自定义协议)
    如果你用了非默认的通信协议,要确保这个映射配置正确,保证advertised.listeners里的协议和listeners里的协议对应上。默认的PLAINTEXT协议不需要额外修改这个配置。

  • 重启Kafka Broker并验证
    改完配置后一定要重启Broker,让新配置生效。重启后可以用Kafka的命令行工具验证元数据里的Broker地址是否正确:

    kafka-topics.sh --describe --bootstrap-server 192.168.95.217:9092 --topic test1
    

    看输出里的Leader和Replicas对应的地址是不是192.168.95.217:9092,而不是localhost。

另外,你的Consumer代码里有个小细节可以优化:kafkaConsumer.seekToBeginning(Collections.emptyList());这句是无效的,因为刚订阅主题时还没有分配分区,传入空列表不会起作用。你已经配置了auto.offset.reset=earliest,这已经能保证从最开始消费,所以这句可以删掉,避免代码混淆。


内容的提问来源于stack exchange,提问作者M. Gopal

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 08:03:51