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

Reactive REST接口无法从外部库返回Flux的问题求助

问题

基于io.projectreactor:reactor-core:3.6.5构建的连接器应用,通过WebSocket订阅字符串数据,映射过滤后返回Flux<Item>。该连接器作为库被Spring Boot REST项目集成,REST控制器代码如下:

@RestController
@RequiredArgsConstructor
public class BrokerApiImpl implements Broker {

    private final BrokerService brokerService;

    @GetMapping(value = "tickerized_orderbook/{trading}/{profit}", produces = MediaType.APPLICATION_JSON_VALUE)
    public Flux<Tick> subscribeTickerizedOrerbook(@PathVariable("trading") String trading,
                                                  @PathVariable("profit") String profit) {
        return brokerService.subscribeTickerizedOrderbook(OrderBookTickSymbol.builder()
                        .currency(CommonCurrencyPair.of(trading, profit))
                        .build())
                .subscribeOn(Schedulers.boundedElastic());
    }
}

该API注册在Eureka Server中,客户端通过WebClient调用并订阅Flux:

@Bean
public List<WebClient> createAllWebClientBrokers(EurekaClient eurekaClient) {
    HttpClient httpClient = HttpClient.create();

    return eurekaClient
            .getApplication("broker")
            .getInstances().stream()
            .map(q -> WebClient
                    .builder()
                    .baseUrl("http://" + q.getIPAddr() + ":" + q.getPort() + "/tickerized_orderbook/{trading}/{profit}")
                    .clientConnector(new ReactorClientHttpConnector(httpClient))
                    .build()
            )
            .toList();
}

客户端订阅代码:

client.get()
        .uri(q -> q.build("btc", "usdt"))
        .retrieve()
        .bodyToFlux(Tick.class)
                .subscribe(t ->
        System.out.println(t));

调用时服务端抛出异常:

java.lang.IllegalStateException: block()/blockFirst()/blockLast() are blocking, which is not supported in thread reactor-http-epoll-2
    at reactor.core.publisher.BlockingSingleSubscriber.blockingGet(BlockingSingleSubscriber.java:86)

已知测试情况:

  • 服务启动时直接订阅该Flux可正常接收数据
  • 返回Flux.fromIterable(...)能正常工作
  • 返回外部库的Flux则报错
  • 将Flux声明为Bean注入控制器返回可成功一次,但无法重复订阅

环境:Spring-Boot parent: 3.2.3,Spring-Cloud dependencies: 2023.0.2

原因分析
  1. 外部库Flux含阻塞操作:外部连接器库返回的Flux内部调用了block()/blockFirst()/blockLast()等阻塞方法,而Spring WebFlux的请求处理线程(如reactor-http-epoll-*)属于非阻塞线程池,Reactor的线程检查机制会直接抛出异常。
  2. Flux复用导致失效:Flux是冷流,声明为单例Bean后首次订阅就会被消耗,后续订阅无法再产生数据,因此只能成功一次。
  3. 线程调度不彻底:控制器中添加的subscribeOn(Schedulers.boundedElastic())仅作用于订阅阶段,若外部库在Flux创建阶段就执行阻塞操作,该调度无法将阻塞逻辑转移到弹性线程池,依然会在请求线程中触发阻塞。
修复方案

方案1:修复外部库的阻塞逻辑

检查外部连接器库的实现,将内部的阻塞方法替换为非阻塞Reactor操作:

  • 把block()替换为flatMap()、switchMap()等链式操作
  • 将所有IO或阻塞逻辑调度到Schedulers.boundedElastic()或Schedulers.io()线程池执行

方案2:隔离阻塞操作(无法修改外部库时)

在服务层将外部库的阻塞逻辑完全隔离到弹性线程池:

@Service
public class BrokerService {
    public Flux<Tick> subscribeTickerizedOrderbook(OrderBookTickSymbol symbol) {
        // 用fromCallable包装可能阻塞的初始化逻辑,提交到弹性线程池
        return Mono.fromCallable(() -> externalLibrary.getFlux(symbol))
                .subscribeOn(Schedulers.boundedElastic())
                .flatMapMany(Function.identity());
    }
}

同时移除控制器中的subscribeOn,避免重复调度:

@GetMapping(value = "tickerized_orderbook/{trading}/{profit}", produces = MediaType.APPLICATION_JSON_VALUE)
public Flux<Tick> subscribeTickerizedOrerbook(@PathVariable("trading") String trading,
                                              @PathVariable("profit") String profit) {
    return brokerService.subscribeTickerizedOrderbook(OrderBookTickSymbol.builder()
                    .currency(CommonCurrencyPair.of(trading, profit))
                    .build());
}

方案3:确保Flux每次请求生成新实例

禁止将外部库返回的Flux声明为单例Bean,确保每次REST请求都创建新的Flux实例:

  • 服务层方法每次调用都返回新的Flux,而非复用已创建的实例
  • 不在Bean初始化阶段创建Flux,而是在请求处理时动态生成

客户端调用优化(可选)

确保客户端WebClient的订阅逻辑隔离到非阻塞线程:

client.get()
        .uri(q -> q.build("btc", "usdt"))
        .retrieve()
        .bodyToFlux(Tick.class)
        .subscribeOn(Schedulers.boundedElastic())
        .subscribe(t -> System.out.println(t));

内容的提问来源于stack exchange,提问作者greengold

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 23:47:46