Java+Kafka+ZooKeeper+Docker环境下会话关闭异常及KafkaProducer问题
嘿,我一眼就看到你代码里的关键问题啦,咱们一步步来解决:
核心问题:参数配置错误
你在初始化KafkaProducer的时候,把ZooKeeper的地址传给了ProducerConfig.BOOTSTRAP_SERVERS_CONFIG,但这个参数需要的是Kafka Broker的地址(格式是host:port,多节点用逗号分隔)!
Kafka Producer从0.10.x版本开始就不需要直接连接ZooKeeper了,它是直接和Broker通信的。当你把ZooKeeper地址填到Broker配置项里,Producer会尝试用Kafka的协议去连接ZooKeeper,自然会连接失败,抛出这个会话关闭的空异常。
看你这段代码:
public KafkaEventPublisher(String zookeeper, Serializer serializer) { Properties props = new Properties(); // 这里是错误根源!把ZooKeeper地址填给了Broker配置 props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, zookeeper); props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class); props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, Stri...
解决步骤
修正Broker地址配置
把BOOTSTRAP_SERVERS_CONFIG的值换成你的Kafka Broker地址,比如本地单节点的话就是localhost:9092,集群环境就是broker1:9092,broker2:9092,broker3:9092。补全Serializer配置
你代码里VALUE_SERIALIZER_CLASS_CONFIG只写了Stri...,得补全完整的Serializer类名。如果是字符串序列化就用StringSerializer.class,如果是自定义序列化,要确保类正确实现了org.apache.kafka.common.serialization.Serializer接口。验证Broker状态与网络连通性
- 确认你的Kafka Broker服务是正常启动运行的
- 检查Producer所在机器能ping通Broker的主机,并且9092端口(默认端口)没有被防火墙拦截
额外排查技巧
如果按上面的步骤改完还是有问题,可以开启Kafka客户端的DEBUG日志,比如在logback或log4j配置里添加:
<logger name="org.apache.kafka" level="DEBUG"/>
这样能看到连接过程中的详细错误信息,更容易定位问题。另外也要确保你的Kafka客户端版本和Broker版本尽量一致,避免版本不兼容导致的奇怪问题。
内容的提问来源于stack exchange,提问作者Dennis van der Veeke

