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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.14 08:45:41