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

基于Spring Cloud Stream的Kafka多主题消费配置方案咨询

Spring Cloud Stream Kafka 多主题消费配置方案

问题背景

在Spring Boot应用中需消费两个Kafka主题:

  • 主题1:所有消息采用统一逻辑处理
  • 主题2:根据消息Header值执行两种不同处理逻辑
    原@StreamListener已废弃,尝试两种方案均遇问题:
  1. 定义两个MessageRoutingCallback Bean时启动失败,提示存在多个同类型Bean
  2. 在配置类中声明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+拆分独立消费者的方式解决冲突,同时满足逻辑拆分需求:

  1. 只定义一个MessageRoutingCallback Bean,根据消息来源主题或Header判断路由目标
  2. 将各处理逻辑封装为独立Consumer Bean,通过Bean名称被路由器调用
  3. 配置绑定多个主题到路由函数,消息经路由后分发到对应消费者

代码实现

统一路由器: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

方案优势

  • 解决多个MessageRoutingCallback Bean的冲突问题
  • 处理逻辑拆分到独立类,代码结构清晰,便于维护
  • 路由逻辑集中统一,便于修改和扩展
  • 配置简洁,通过单绑定实现多主题监听

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 03:25:04