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
相关产品推荐
相关产品推荐

