无法配置Spring Cloud Stream Kafka Binder连接Broker测试失败求助
问题描述
使用Stream Bridge发送Kafka消息,代码为streamBridge.send("alm-foo-dev", "kafka", message),测试时无法连接Kafka Broker。日志显示无法连接localhost:9952,且抛出java.lang.IllegalStateException: No records found for topic异常。
相关配置如下:
EmbeddedKafka配置
@EmbeddedKafka( brokerProperties = {"listeners=PLAINTEXT://localhost:9952", "port=9952"}, topics = {"alm-jira-dev"}, partitions = 1 )
测试用application.yml配置
server: port: 29191 logging: level: root: ERROR org: springframework.integration: DEBUG springframework.cloud.stream: DEBUG springframework.boot.autoconfigure.mongo: WARN com.digite.cloud: DEBUG spring: application: name: howler main: banner-mode: off mongodb: embedded: version: 3.4.6 data: mongodb: port: 29129 host: localhost database: howler_db cloud: stream: kafka: default: producer: useTopicHeader: true binder: defaultBrokerPort: 9952 autoCreateTopics: false producerProperties: key.serializer: org.apache.kafka.common.serialization.StringSerializer value.serializer: org.springframework.kafka.support.serializer.JsonSerializer max.block.ms: 100
解决方案
- Topic名称不匹配:EmbeddedKafka预创建的topic是
alm-jira-dev,但发送消息指定的是alm-foo-dev,这会直接导致找不到对应topic。解决方式二选一:将发送代码中的topic改为alm-jira-dev,或者在EmbeddedKafka的topics参数中添加alm-foo-dev。 - Stream Bridge调用参数错误:
streamBridge.send的第一个参数应为绑定名称或目标topic,第二个参数"kafka"是多余的,会被识别为绑定名称的一部分,导致无法匹配正确的Kafka Binder。正确的调用方式为streamBridge.send("alm-foo-dev", message),或者显式指定目标类型:streamBridge.send("destination:alm-foo-dev", message)。 - 完善Kafka Binder地址配置:仅配置
defaultBrokerPort可能不足以让Binder正确识别Broker地址,需补充brokers配置:spring: cloud: stream: kafka: binder: brokers: localhost:9952 defaultBrokerPort: 9952 # 保留其他原有配置 - 自动创建Topic配置冲突:当前设置
autoCreateTopics: false,但如果发送的topic未在EmbeddedKafka中预先创建,会导致找不到topic。可选择开启autoCreateTopics: true,或确保EmbeddedKafka的topics包含所有需要发送的topic。 - 调整连接超时时间:
max.block.ms: 100设置过小,可能在Broker连接建立完成前就触发超时,建议调大至5000或使用默认值,避免因超时导致的连接失败。
内容的提问来源于stack exchange,提问作者Anadi Misra
相关产品推荐
相关产品推荐

