Spring Cloud Function组合函数中首个函数被重复调用的原因排查
我来帮你分析下这个问题——你遇到的首个函数(uppercase)日志打印两次,但下游的log Consumer只执行一次的情况,其实和Reactive流的订阅特性以及Spring Cloud Function处理组合函数的方式有关,咱们一步步拆解:
1. 核心原因:Reactive冷流的多订阅触发
你的uppercase函数返回的是Flux<Message<String>>,这是一个冷流(Cold Stream)——冷流的核心特点是:每被订阅一次,就会从头完整执行一遍流的处理逻辑。
当你把uppercase和log组合成uppercase|log时,Spring Cloud Function在处理这个组合链的过程中,会对uppercase输出的流进行两次订阅:
- 第一次订阅是为了把流的内容传递给下游的
logConsumer; - 第二次订阅可能来自WebFlux框架本身(比如需要确认请求的完成状态、响应的
onComplete信号),或是Spring Cloud Function的Web绑定模块在适配Reactive函数时的额外校验逻辑。
而下游的log Consumer只执行一次,是因为它是一个命令式(Imperative)的Consumer,Spring在适配Reactive流到命令式Consumer时,会做一层单订阅适配,确保Consumer只会被触发一次,不会跟着流的多订阅重复执行。
2. 代码层面的验证细节
你可以从现有现象里找到佐证:
- 测试用例通过(
handler.handleMessage只被调用一次):因为TestRestTemplate是阻塞式客户端,测试时的适配层会确保流只被有效订阅一次来触发Consumer; uppercase的日志打印两次:正好对应了流被两次订阅的情况,每次订阅都会执行一遍flux.map里的逻辑。
3. 解决办法:将冷流转为热流,避免重复执行
要让uppercase的逻辑只执行一次,你可以把冷流转换成热流(Hot Stream),让多次订阅共享同一份流的处理结果。修改uppercase函数,添加publish().refCount()操作符:
@Bean Function<Flux<Message<String>>, Flux<Message<String>>> uppercase() { return flux -> flux.map(message -> { final var payload = message.getPayload(); LOGGER.info("Will transform this payload to upper case: {}", payload); return MessageBuilder .withPayload(payload.toUpperCase()) .copyHeaders(message.getHeaders()) .build(); }) // 将冷流转为热流,多订阅共享同一处理结果 .publish() .refCount(1); // 至少1个订阅者时激活流,无订阅者时自动销毁 }
publish().refCount(1)会把冷流包装成可连接流,第一个订阅者出现时启动流的处理,后续订阅者直接复用已有的流数据,这样无论被订阅多少次,map里的转换逻辑只会执行一次。
4. 替代方案:改用命令式函数(如果业务允许)
如果你的业务场景不需要Reactive特性,也可以直接把uppercase改成命令式的Function,这样Spring Cloud Function会用命令式的方式处理组合链,从根源上避免多订阅问题:
@Bean Function<Message<String>, Message<String>> uppercase() { return message -> { final var payload = message.getPayload(); LOGGER.info("Will transform this payload to upper case: {}", payload); return MessageBuilder .withPayload(payload.toUpperCase()) .copyHeaders(message.getHeaders()) .build(); }; }
内容来源于stack exchange

