Apache Camel Reactive Stream报无活动订阅异常
异常产生原因
- 核心是执行时序违反Reactive Streams规范:
sayHi()方法中先调用template.asyncSendBody()触发消息发送,消息路由到reactive-streams:greet端点时,后续Mono.from(crss.fromStream(...))对应的订阅还未建立。Reactive Streams组件要求发布者推送消息时必须存在已激活的订阅者,否则直接抛出java.lang.IllegalStateException: The stream has no active subscriptions。 - 原有写法中消息发送、流订阅是两个完全独立的操作,没有绑定时序保证:
Mono是惰性执行的,只有当WebFlux作为下游真正订阅返回的Mono时订阅才会生效,而asyncSendBody在方法进入时就立刻触发异步发送,两者存在明显的时间差,必然出现消息先到、订阅后建的问题。 - 额外隐患:类同时作为
RouteBuilder和@RestController使用时,缺少@Component注解会导致Camel无法自动识别路由定义,不过这不是本次异常的直接触发原因。
修复方案
核心原则是先完成订阅准备,再触发消息发送,将两个操作绑定到同一个响应式执行链路中,利用响应式流的回调保证时序:
- 先构造流的
Mono实例,加上.next()取单值适配Mono的单值语义 - 将消息发送逻辑放到
doOnSubscribe回调中,该回调只有当下游真正完成订阅、开始请求数据时才会执行,从根源避免无订阅发消息的问题 - 给类补充
@Component注解保证Camel路由能被正常加载
修复后的核心代码如下:
package com.manning.camel.reactive; import org.apache.camel.ProducerTemplate; import org.apache.camel.builder.RouteBuilder; import org.apache.camel.component.reactive.streams.api.CamelReactiveStreamsService; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Component; import org.springframework.web.bind.annotation.GetMapping; import org.springframework.web.bind.annotation.RestController; import reactor.core.publisher.Mono; @Component // 补充注解保证路由被Camel扫描到 @RestController public class MySpringBootRouter extends RouteBuilder { @Autowired private ProducerTemplate template; @Autowired private CamelReactiveStreamsService crss; @GetMapping public Mono<String> sayHi() { // 先构造流接收逻辑,不立即触发执行 Mono<String> greetStream = Mono.from(crss.fromStream("greet", String.class)) .next(); // 取流中第一个元素,适配Mono单值特性 // 订阅激活后再触发消息发送 return greetStream.doOnSubscribe(subscription -> template.asyncSendBody("direct:works", "Hi") ); } @Override public void configure() { from("direct:works") .log("Fired") .to("reactive-streams:greet"); } }
如果是业务场景需要每次请求独立收发消息,更推荐直接使用Camel的to("reactor:greet")专属端点配合Reactor,写法会更简洁,不需要手动管理订阅时序。
内容的提问来源于stack exchange,提问作者user902870
相关产品推荐
相关产品推荐

