基于Spring Cloud Stream的Kafka多主题消费配置方案咨询
Spring Cloud Stream Kafka 多主题消费配置方案
问题背景
在Spring Boot应用中需消费两个Kafka主题:
- 主题1:所有消息采用统一逻辑处理
- 主题2:根据消息Header值执行两种不同处理逻辑
原@StreamListener已废弃,尝试两种方案均遇问题:
- 定义两个
MessageRoutingCallbackBean时启动失败,提示存在多个同类型Bean - 在配置类中声明Consumer Bean,但不想将所有处理逻辑集中在此,期望更清晰的拆分方案
启动报错信息
*************************** APPLICATION FAILED TO START *************************** Description: Parameter 3 of method functionRouter in org.springframework.cloud.function.context.config.ContextFunctionCatalogAutoConfiguration required a single bean, but 2 were found: - firstMessageRouter: defined in file [C:\Workspace-IntelliJ\TEST\test-kafka-cloud-stream\target\classes\org\example\router\FirstMessageRouter.class] - secondMessageRouter: defined in file [C:\Workspace-IntelliJ\TEST\test-kafka-cloud-stream\target\classes\org\example\router\SecondMessageRouter.class] Action: Consider marking one of the beans as @Primary, updating the consumer to accept multiple beans, or using @Qualifier to identify the bean that should be consumed
尝试过的代码
第一种尝试:两个MessageRoutingCallback Bean
FirstMessageRouter.java
@Component @Slf4j public class FirstMessageRouter implements MessageRoutingCallback { @Override public FunctionRoutingResult routingResult(Message<?> message) { log.info("First router"); log.info("Headers: " + message.getHeaders()); log.info("Payload: " + message.getPayload()); return new MessageRoutingCallback.FunctionRoutingResult("firstconsumer"); } }
SecondMessageRouter.java
@Component @Slf4j public class SecondMessageRouter implements MessageRoutingCallback { @Override public FunctionRoutingResult routingResult(Message<?> message) { log.info("Second router"); log.info("Headers: " + message.getHeaders()); log.info("Payload: " + message.getPayload()); return new FunctionRoutingResult("second-consumer"); } }
FirstConsumer.java
@Slf4j @Component public class FirstConsumer implements Consumer<Message<?>> { @Override public void accept(Message<?> message) { log.info("Received message in FirstConsumer with payload: " + message.getPayload() + " and headers: " + message.getHeaders()); } }
SecondConsumer.java
@Slf4j @Component public class SecondConsumer implements Consumer<Message<?>> { @Override public void accept(Message<?> message) { log.info("Received message in SecondConsumer with payload: " + message.getPayload() + " and headers: " + message.getHeaders()); } }
application.yml(第一种尝试)
spring: cloud: stream: bindings: firstMessageRouter-in-0: destination: topic-1-input group: group-3 secondMessageRouter-in-0: destination: topic-2-input group: group-4 kafka: streams: defaultBrokerPort: 9092 brokers: localhost
第二种尝试:配置类中声明Consumer Bean
KafkaConfiguration.java
@Configuration @Slf4j public class KafkaConfiguration { @Bean("first-process") public Consumer<Message<?>> process() { return message -> { log.info("Received headers: " + message.getHeaders()); log.info("Received payload: " + message.getPayload()); }; } @Bean("second-process") Consumer<String> topic1() { return str -> { log.info("Received message in topic1: " + str); }; } }
application.yml(第二种尝试)
spring: cloud: stream: function: definition: first-process; second-process bindings: first-process-in-0: destination: topic.input group: group-1 second-process-in-0: destination: topic-1 group: group-2 kafka: streams: defaultBrokerPort: 9092 brokers: localhost
解决方案
方案思路
通过单统一路由Bean+拆分独立消费者的方式解决冲突,同时满足逻辑拆分需求:
- 只定义一个
MessageRoutingCallbackBean,根据消息来源主题或Header判断路由目标 - 将各处理逻辑封装为独立
ConsumerBean,通过Bean名称被路由器调用 - 配置绑定多个主题到路由函数,消息经路由后分发到对应消费者
代码实现
统一路由器:TopicMessageRouter.java
@Component @Slf4j public class TopicMessageRouter implements MessageRoutingCallback { @Override public FunctionRoutingResult routingResult(Message<?> message) { String topic = message.getHeaders().get(KafkaHeaders.RECEIVED_TOPIC, String.class); log.info("Received message from topic: {}", topic); // 主题1统一路由到firstConsumer if ("topic-1-input".equals(topic)) { return new FunctionRoutingResult("firstConsumer"); } // 主题2根据Header值路由到不同消费者 else if ("topic-2-input".equals(topic)) { String processType = message.getHeaders().get("process-type", String.class); return "typeA".equals(processType) ? new FunctionRoutingResult("secondConsumerTypeA") : new FunctionRoutingResult("secondConsumerTypeB"); } // 默认路由(可选) return new FunctionRoutingResult("defaultConsumer"); } }
拆分的消费者类
FirstConsumer.java(主题1统一处理)
@Slf4j @Component("firstConsumer") public class FirstConsumer implements Consumer<Message<?>> { @Override public void accept(Message<?> message) { log.info("Topic1处理逻辑 - Payload: {}, Headers: {}", message.getPayload(), message.getHeaders()); } }
SecondConsumerTypeA.java(主题2 Header=typeA处理)
@Slf4j @Component("secondConsumerTypeA") public class SecondConsumerTypeA implements Consumer<Message<?>> { @Override public void accept(Message<?> message) { log.info("Topic2 TypeA处理逻辑 - Payload: {}, Headers: {}", message.getPayload(), message.getHeaders()); } }
SecondConsumerTypeB.java(主题2 Header=typeB处理)
@Slf4j @Component("secondConsumerTypeB") public class SecondConsumerTypeB implements Consumer<Message<?>> { @Override public void accept(Message<?> message) { log.info("Topic2 TypeB处理逻辑 - Payload: {}, Headers: {}", message.getPayload(), message.getHeaders()); } }
DefaultConsumer.java(可选默认处理)
@Slf4j @Component("defaultConsumer") public class DefaultConsumer implements Consumer<Message<?>> { @Override public void accept(Message<?> message) { log.info("默认处理逻辑 - 未知主题/Header的消息: {}", message.getPayload()); } }
application.yml配置
spring: cloud: stream: function: definition: functionRouter bindings: functionRouter-in-0: # 绑定多个主题,逗号分隔 destination: topic-1-input,topic-2-input group: kafka-consumer-group consumer: # 启用多主题监听 multi: true kafka: binder: brokers: localhost:9092 bindings: functionRouter-in-0: consumer: # 启用Header传递,确保能获取Kafka主题和自定义Header header-mode: headers
方案优势
- 解决多个
MessageRoutingCallbackBean的冲突问题 - 处理逻辑拆分到独立类,代码结构清晰,便于维护
- 路由逻辑集中统一,便于修改和扩展
- 配置简洁,通过单绑定实现多主题监听
内容的提问来源于stack exchange,提问作者Alexxxx
相关产品推荐
相关产品推荐

