Spring Boot WebFlux中Flux多订阅者分执行上下文执行问题问询
问题拆解与解决方案
先帮你理清两个核心问题的本质,再给出针对性的解决办法:
1. 为什么两个订阅者拿到不同的事件值?
你用的Flux.interval()是冷流——响应式流里的冷流每次被订阅都会从头生成新的序列。你的代码里,events.subscribe(...)是第一个订阅,webSocketSession.send(...)内部会触发第二个订阅,这两个订阅各自触发了map(i -> new MyEvent(随机值))的执行,所以生成的随机值完全不同,日志里的数值自然不一致。
你后来用share()解决了这个问题,是因为share()把冷流转换成了热流——热流只会生成一次序列,所有订阅者共享同一个事件流,所以事件值就一致了。
2. 为什么用share()后会出现阻塞?
问题出在Notifier里的RestTemplate:它是同步阻塞的HTTP客户端,当它在Flux默认的响应式线程(比如parallel线程池)里执行时,会占用这个线程,导致WebSocket的消息发送被阻塞。而且share()后的流默认会在同一个调度器上运行,两个订阅者的逻辑互相抢占线程资源,自然影响了WebSocket的实时性。
解决方案:让两个订阅者在独立执行上下文运行
方案1:给阻塞操作分配专门的线程池(适合保留RestTemplate的场景)
第一步:配置自定义阻塞线程池
创建一个配置类,专门生成用于处理阻塞操作的调度器:
@Configuration public class SchedulerConfig { @Bean public Scheduler blockingNotifierScheduler() { // newBoundedElastic是专门为阻塞操作设计的调度器,自动管理线程数 return Schedulers.newBoundedElastic( 10, // 核心线程数 100, // 最大线程数 "blocking-notifier-pool" // 线程池名称,方便日志排查 ); } }
第二步:修改订阅逻辑,切换执行上下文
在MyWebSocketHandler里,把通知器的逻辑切换到自定义的阻塞线程池,和WebSocket的线程彻底隔离:
public class MyWebSocketHandler implements WebSocketHandler { @Autowired private MySource mySource; @Autowired private Notifier notifier; @Autowired private Scheduler blockingNotifierScheduler; @Override public Mono<Void> handle(WebSocketSession webSocketSession) { Flux<MyEvent> sharedEvents = mySource.events(); // 把通知器的逻辑切换到阻塞线程池执行,不占用WebSocket的响应式线程 sharedEvents .publishOn(blockingNotifierScheduler) // 指定后续操作在自定义线程池执行 .subscribe(event -> notifier.sendNotification(event)); // WebSocket的逻辑保持在响应式默认线程,不受阻塞操作影响 return webSocketSession.send( sharedEvents .map(this::toJson) .map(webSocketSession::textMessage) ); } private String toJson(MyEvent event) { log.info("websocket toJson " + event.getValue()); // 这里用ObjectMapper序列化,示例代码 try { return new ObjectMapper().writeValueAsString(event); } catch (JsonProcessingException e) { throw new RuntimeException("Failed to serialize event", e); } } }
方案2:替换RestTemplate为WebClient(更贴合WebFlux的非阻塞理念)
如果想彻底避免阻塞线程的问题,建议把同步的RestTemplate换成WebFlux自带的非阻塞WebClient,这样整个通知逻辑都是响应式的,不需要额外的线程池:
修改Notifier类:
public class Notifier { private final WebClient webClient; // 构造注入WebClient.Builder,推荐这种方式配置WebClient public Notifier(WebClient.Builder webClientBuilder) { this.webClient = webClientBuilder .baseUrl("https://your-notification-api-endpoint.com") // 替换成你的通知API地址 .build(); } public Mono<Void> sendNotification(MyEvent event) { log.info("notifier sendNotification " + event.getValue()); if (event.getValue() > 4) { // 示例条件,替换成你的校验逻辑 return webClient.post() .uri("/send-notification") // 替换成你的API路径 .bodyValue(event) .retrieve() .bodyToMono(Void.class) .doOnError(ex -> log.error("Notification failed for event: {}", event, ex)); } return Mono.empty(); // 条件不满足时返回空Mono,无操作 } }
修改MyWebSocketHandler的订阅逻辑:
@Override public Mono<Void> handle(WebSocketSession webSocketSession) { Flux<MyEvent> sharedEvents = mySource.events(); // 响应式订阅通知逻辑,WebClient本身是非阻塞的,不会占用线程 sharedEvents .flatMap(event -> notifier.sendNotification(event)) .subscribe(); // 订阅触发执行,忽略结果(如果需要处理失败可以加doOnError) // WebSocket的发送逻辑保持不变 return webSocketSession.send( sharedEvents .map(this::toJson) .map(webSocketSession::textMessage) ); }
效果验证
做完上述修改后:
- 两个订阅者会拿到完全一致的事件值(因为用了
share()的热流) - WebSocket的消息发送不会被REST通知阻塞,两者在独立的线程/执行上下文运行,互不影响
内容的提问来源于stack exchange,提问作者vortex.alex
相关产品推荐
相关产品推荐

