You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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:所有主题无注册消费者

问题分析与解决方案

核心原因

  1. 函数输入绑定未激活:Spring Cloud Stream默认只会为函数绑定第一个输入(in-0),第二个输入(in-1)不会自动生效,需显式声明函数支持多输入。
  2. 自定义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\";"

验证方法

  1. 重启应用,查看日志中是否有kafka-two消费者连接成功的记录
  2. 通过Kafka管理工具检查kafka-two的order-processed主题是否有spot消费者组注册
  3. 向kafka-two的order-processed主题发送测试消息,确认orderProcessedListener能正常接收

内容的提问来源于stack exchange,提问作者Jonathan Henrique

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.18 11:17:12