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

