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

Spring Cloud Stream多Kafka集群配置报错,如何调整?

多Kafka集群Spring Cloud Stream配置问题解决

问题背景

我需要在项目中连接两个不同的Kafka Broker,当前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
        orderProcessedListener-in-0:
          destination: order-processed
          group: spot
          binder: kafka-two
      kafka:
        binder:
          auto-create-topics: true
        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
      binders:
        kafka-one:
          type: kafka
          environment:
            spring:
              cloud:
                stream:
                  kafka:
                    binder:
                      brokers: localhost:9092
        kafka-two:
          type: kafka
          environment:
            spring:
              cloud:
                stream:
                  kafka:
                    binder:
                      brokers: localhost:9093

运行应用时出现以下连接超时报错:

2024-03-05T23:35:48.473-03:00  INFO 25569 --- [| adminclient-1] org.apache.kafka.clients.NetworkClient   : [AdminClient clientId=adminclient-1] Cancelled in-flight API_VERSIONS request with correlation id 31 due to node 1001 being disconnected (elapsed time since creation: 4ms, elapsed time since send: 4ms, request timeout: 3600000ms)
2024-03-05T23:35:49.595-03:00  INFO 25569 --- [| adminclient-1] org.apache.kafka.clients.NetworkClient   : [AdminClient clientId=adminclient-1] Node 1001 disconnected.
2024-03-05T23:35:49.595-03:00  INFO 25569 --- [| adminclient-1] org.apache.kafka.clients.NetworkClient   : [AdminClient clientId=adminclient-1] Cancelled in-flight API_VERSIONS request with correlation id 32 due to node 1001 being disconnected (elapsed time since creation: 5ms, elapsed time since send: 5ms, request timeout: 3600000ms)
2024-03-05T23:35:50.727-03:00  INFO 25569 --- [| adminclient-1] org.apache.kafka.clients.NetworkClient   : [AdminClient clientId=adminclient-1] Node 1001 disconnected.
2024-03-05T23:35:50.728-03:00  INFO 25569 --- [| adminclient-1] org.apache.kafka.clients.NetworkClient   : [AdminClient clientId=adminclient-1] Cancelled in-flight API_VERSIONS request with correlation id 33 due to node 1001 being disconnected (elapsed time since creation: 4ms, elapsed time since send: 4ms, request timeout: 3600000ms)
2024-03-05T23:35:51.086-03:00  INFO 25569 --- [| adminclient-1] o.a.k.c.a.i.AdminMetadataManager         : [AdminClient clientId=adminclient-1] Metadata update failed

org.apache.kafka.common.errors.TimeoutException: Timed out waiting for a node assignment. Call: fetchMetadata

我的需求是将Kafka主题拆分到两个集群:kafka-one包含order-created和order-created-dlq,kafka-two包含order-processed和order-processed-dlq。使用技术版本为Spring Boot 3.2.3、Spring Cloud 2023.0.0,两个Kafka集群通过Docker容器在开发环境正常运行,分别暴露9092和9093端口。


问题根源

全局层级的spring.cloud.stream.kafka配置会被所有Binder继承,而DLQ相关配置绑定在全局节点下,导致AdminClient错误地用非目标集群连接去创建DLQ主题,最终引发连接超时。核心问题是没有将每个Binder的专属配置(包括DLQ)完全隔离到对应Binder的环境中。

调整后的配置

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
        orderProcessedListener-in-0:
          destination: order-processed
          group: spot
          binder: kafka-two
      binders:
        kafka-one:
          type: kafka
          environment:
            spring:
              cloud:
                stream:
                  kafka:
                    binder:
                      brokers: localhost:9092
                      auto-create-topics: true
                    bindings:
                      orderCreatedListener-in-0:
                        consumer:
                          enableDlq: true
                          dlqName: order-created-dlq
                          autoCommitOnError: true
                          autoCommitOffset: true
        kafka-two:
          type: kafka
          environment:
            spring:
              cloud:
                stream:
                  kafka:
                    binder:
                      brokers: localhost:9093
                      auto-create-topics: true
                    bindings:
                      orderProcessedListener-in-0:
                        consumer:
                          enableDlq: true
                          dlqName: order-processed-dlq
                          autoCommitOnError: true
                          autoCommitOffset: true

关键调整说明

  • 移除全局的spring.cloud.stream.kafka节点,将auto-create-topics分别放到kafka-one和kafka-two的Binder环境配置中,避免全局配置干扰
  • 把原来全局kafka.bindings下的消费者DLQ配置,分别迁移到对应Binder的environment.spring.cloud.stream.kafka.bindings下,确保每个DLQ主题由指定集群的AdminClient创建
  • 每个Binder拥有独立的完整配置,彻底避免跨集群的元数据请求冲突

内容的提问来源于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.28 19:02:00