Spring Cloud Stream单消费者多Kafka集群多输入绑定配置问题
问题描述
需要连接两个不同的Kafka Broker,分别消费每个集群上order-processed Topic的事件,但运行应用时,Kafka One的消费者绑定orderProcessedListener-in-0正常工作,Kafka Two的orderProcessedListener-in-1绑定无法工作。
使用环境
- Spring Boot 3.3.0
- Spring Cloud 2023.0.1
- 开发环境:两个Kafka集群通过Docker容器运行,分别暴露9092和9093端口
当前配置(application.yaml)
spring: cloud: function: definition: orderCreatedListener;orderProcessedListener stream: bindings: orderCreatedProducer-out-0: destination: order-created binder: kafka-one orderCreatedListener-in-0: destination: order-created group: spot binder: kafka-one orderCreatedListener-out-0: destination: order-processed binder: kafka-two # 测试时手动切换该binder,orderProcessedListener-in-1无专属生产者 orderProcessedListener-in-0: # 从KAFKA ONE消费 destination: order-processed group: spot binder: kafka-one orderProcessedListener-in-1: # 从KAFKA TWO消费 destination: order-processed group: spot binder: kafka-two kafka: binder: auto-create-topics: true configuration: security: protocol: SASL_PLAINTEXT sasl: mechanism: PLAIN bindings: orderCreatedListener-in-0: consumer: enableDlq: true dlqName: order-created-dlq autoCommitOnError: true autoCommitOffset: true orderProcessedListener-in-0: consumer: enableDlq: true dlqName: order-processed-dlq autoCommitOnError: true autoCommitOffset: true orderProcessedListener-in-1: consumer: enableDlq: true dlqName: order-processed-dlq autoCommitOnError: true autoCommitOffset: true binders: kafka-one: type: kafka environment: spring: cloud: stream: kafka: binder: brokers: localhost:9092 configuration: sasl: jaas: config: "org.apache.kafka.common.security.plain.PlainLoginModule required username=\"spot\" password=\"spot\";" kafka-two: type: kafka environment: spring: cloud: stream: kafka: binder: brokers: localhost:9093 configuration: sasl: jaas: config: "org.apache.kafka.common.security.plain.PlainLoginModule required username=\"spot\" password=\"spot\";"
集群状态
- Kafka One:所有主题已注册消费者,状态正常
- Kafka Two:所有主题无注册消费者
问题分析与解决方案
核心原因
- 函数输入绑定未激活:Spring Cloud Stream默认只会为函数绑定第一个输入(
in-0),第二个输入(in-1)不会自动生效,需显式声明函数支持多输入。 - 自定义Binder配置缺失:全局Kafka安全配置不会继承到自定义Binder,kafka-two缺少完整的SASL协议配置,导致无法建立有效连接。
调整步骤
1. 激活函数的多输入绑定
修改spring.cloud.function配置,明确指定orderProcessedListener支持2个输入:
spring: cloud: function: definition: orderCreatedListener;orderProcessedListener#2
#2表示该函数接受2个输入参数,会自动绑定in-0和in-1两个通道。
2. 完善kafka-two的安全配置
在kafka-two的环境配置中补充完整的安全协议和SASL机制,确保和kafka-one的配置对齐:
kafka-two: type: kafka environment: spring: cloud: stream: kafka: binder: brokers: localhost:9093 auto-create-topics: true configuration: security: protocol: SASL_PLAINTEXT sasl: mechanism: PLAIN jaas: config: "org.apache.kafka.common.security.plain.PlainLoginModule required username=\"spot\" password=\"spot\";"
3. 完整调整后的配置示例
spring: cloud: function: definition: orderCreatedListener;orderProcessedListener#2 stream: bindings: orderCreatedProducer-out-0: destination: order-created binder: kafka-one orderCreatedListener-in-0: destination: order-created group: spot binder: kafka-one orderCreatedListener-out-0: destination: order-processed binder: kafka-two orderProcessedListener-in-0: destination: order-processed group: spot binder: kafka-one orderProcessedListener-in-1: destination: order-processed group: spot binder: kafka-two kafka: binder: auto-create-topics: true configuration: security: protocol: SASL_PLAINTEXT sasl: mechanism: PLAIN bindings: orderCreatedListener-in-0: consumer: enableDlq: true dlqName: order-created-dlq autoCommitOnError: true autoCommitOffset: true orderProcessedListener-in-0: consumer: enableDlq: true dlqName: order-processed-dlq autoCommitOnError: true autoCommitOffset: true orderProcessedListener-in-1: consumer: enableDlq: true dlqName: order-processed-dlq autoCommitOnError: true autoCommitOffset: true binders: kafka-one: type: kafka environment: spring: cloud: stream: kafka: binder: brokers: localhost:9092 auto-create-topics: true configuration: security: protocol: SASL_PLAINTEXT sasl: mechanism: PLAIN jaas: config: "org.apache.kafka.common.security.plain.PlainLoginModule required username=\"spot\" password=\"spot\";" kafka-two: type: kafka environment: spring: cloud: stream: kafka: binder: brokers: localhost:9093 auto-create-topics: true configuration: security: protocol: SASL_PLAINTEXT sasl: mechanism: PLAIN jaas: config: "org.apache.kafka.common.security.plain.PlainLoginModule required username=\"spot\" password=\"spot\";"
验证方法
- 重启应用,查看日志中是否有kafka-two消费者连接成功的记录
- 通过Kafka管理工具检查kafka-two的
order-processed主题是否有spot消费者组注册 - 向kafka-two的
order-processed主题发送测试消息,确认orderProcessedListener能正常接收
内容的提问来源于stack exchange,提问作者Jonathan Henrique
相关产品推荐
相关产品推荐

