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

Spring Cloud Stream函数式消费如何实现基于内容的路由

Spring Cloud Stream 函数式编程模式下实现内容路由的方案

你使用的Spring Boot 2.3.12.RELEASE + Hoxton.SR11版本组合对应的Spring Cloud Stream为3.0.x版本,可通过以下两种方式实现和原StreamListener condition效果一致的内容路由:

方案一:单Consumer内部手动路由(最简洁,无需额外配置)

直接将Consumer的入参改为Spring的Message类型,获取消息头后做逻辑判断,分发到不同的处理方法即可,代码示例如下:

import org.springframework.messaging.Message;

@Bean
public Consumer<Message<Object>> fileRequestConsumer() {
    return message -> {
        String sagaType = message.getHeaders().get("saga_request", String.class);
        try {
            if ("FILE_SUBMIT".equals(sagaType)) {
                handleSubmitFile((FileSubmitRequest) message.getPayload());
            } else if ("FILE_CANCEL".equals(sagaType)) {
                handleCancelFile((FileCancelRequest) message.getPayload());
            } else {
                // 可自定义未匹配消息的处理逻辑,比如输出日志、投递到死信队列
                LOGGER.warn("未匹配的saga_request类型:{}", sagaType);
            }
        } catch (Exception e) {
            LOGGER.error("消息处理异常", e);
        }
    };
}

private void handleSubmitFile(FileSubmitRequest request) {
    // 原有文件提交的业务逻辑
}

private void handleCancelFile(FileCancelRequest request) {
    // 原有文件取消的业务逻辑
}

方案二:使用框架自带的路由表达式配置(和原StreamListener用法最接近)

Spring Cloud Stream 3.x版本内置了函数路由能力,你只需要分别定义不同逻辑的Consumer Bean,再通过配置路由表达式即可实现框架自动路由,无需自己写判断逻辑。

第一步:定义独立的处理函数

@Bean
public Consumer<FileSubmitRequest> handleSubmitFile() {
    return request -> {
        // 文件提交业务逻辑
    };
}

@Bean
public Consumer<FileCancelRequest> handleCancelFile() {
    return request -> {
        // 文件取消业务逻辑
    };
}

// 可选:定义未匹配消息的默认处理函数
@Bean
public Consumer<Message<Object>> defaultConsumer() {
    return message -> LOGGER.warn("未匹配的消息:{}", message);
}

第二步:添加配置

在application.yml中添加路由相关配置:

spring:
  cloud:
    stream:
      bindings:
        # routingFunction-in-0为框架路由功能默认的输入通道绑定名,绑定到你的原FILE_REQUEST_IN对应的Topic
        routingFunction-in-0:
          destination: 你的业务Topic名称
          group: 你的消费组名称
      function:
        # 路由表达式,根据saga_request头匹配对应的处理函数Bean名
        routing-expression: "headers['saga_request'] == 'FILE_SUBMIT' ? 'handleSubmitFile' : headers['saga_request'] == 'FILE_CANCEL' ? 'handleCancelFile' : 'defaultConsumer'"

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.29 23:57:02