使用@EmbeddedKafka时为何必须配置brokerProperties?
这个问题其实是旧版本spring-kafka和Spring Boot组合使用时的常见小坑,我来给你拆解清楚~
首先得搞懂两个关键角色的默认行为:
- EmbeddedKafka的默认配置:在你使用的spring-kafka 2.9.9(对应Spring Boot 2.7.x)版本中,如果你不指定
listeners参数,它会随机挑选一个可用端口,并且绑定到0.0.0.0(允许所有网卡访问),但它不会主动把自己的实际运行地址同步给应用内的KafkaTemplate。 - KafkaTemplate的默认配置:Spring Boot的自动配置逻辑会默认把KafkaTemplate的
bootstrap.servers设为localhost:9092——这是Kafka官方的默认端口。
这就产生了矛盾:EmbeddedKafka在随机端口上运行,而你的KafkaTemplate却一个劲往localhost:9092发送请求,自然连不上目标Broker,也就出现了你看到的Connection to node -1 (localhost/127.0.0.1:9092) could not be established报错。
那为什么加上listeners=PLAINTEXT://localhost:9092就正常了?因为这行配置强制EmbeddedKafka绑定到localhost:9092这个固定端口,正好和KafkaTemplate默认的连接地址对齐,两者通信地址一致,自然就能正常交互了。
可能你会疑惑:EmbeddedKafka不该自动把自己的地址注入到应用配置里吗?没错,在spring-kafka 3.0+(对应Spring Boot 3.x及以上)的新版本中,这个逻辑已经被优化了——EmbeddedKafka启动后会自动把自身的实际端口和地址填充到Spring环境变量中,KafkaTemplate会自动读取这些值,完全不需要手动配置listeners。但你用的是2.9.9这个旧版本,还没有这个自动适配的功能,所以必须手动指定listeners来对齐端口。
另外还有一种更灵活的解决方案,不用绑定固定端口:你可以通过EmbeddedKafkaBroker获取它的实际运行地址,手动配置KafkaTemplate的连接参数,示例代码如下:
@Autowired private EmbeddedKafkaBroker embeddedKafkaBroker; @Bean public KafkaTemplate<String, String> kafkaTemplate() { Map<String, Object> configs = new HashMap<>(KafkaTestUtils.producerProps(embeddedKafkaBroker)); return new KafkaTemplate<>(new DefaultKafkaProducerFactory<>(configs)); }
这种方式下,不管EmbeddedKafka使用哪个随机端口,KafkaTemplate都能正确连接,不用写死9092端口,不过需要你手动创建KafkaTemplate Bean,而非使用Spring自动配置的实例。
总结一下:你需要添加那行配置的核心原因是旧版本spring-kafka的EmbeddedKafka不会自动同步运行地址给KafkaTemplate,导致两者连接地址不匹配,手动指定listeners让EmbeddedKafka绑定到KafkaTemplate默认的9092端口,就能解决这个问题。
备注:内容来源于stack exchange,提问作者fml2

