远程运行Kafka Java 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

