Spring 5中FluxSink调用next不实时发数据的问题咨询
先解释第一个现象的原因
在你第一个版本的代码里,Flux.create的回调逻辑是直接跑在**Netty的EventLoop线程(也就是处理HTTP请求的主线程)**上的。Reactor框架对EventLoop线程有个核心要求:绝对不能阻塞它!
这个EventLoop线程身兼数职:既要处理请求的接收,也要负责把响应数据推送给客户端。当你在这个线程里调用Thread.sleep(2000)的时候,相当于把它死死占住了——它根本没时间去处理每次fluxSink.next()之后的响应推送工作,只能等整个for循环跑完、fluxSink.complete()执行完毕,才会一次性把攒下来的所有数据打包发给客户端。这就导致你看不到实时推送的效果。
而第二个版本把生成数据的逻辑放到了独立线程里,这样EventLoop线程立刻就被释放了。它可以在每次fluxSink.next()调用时,及时把新生成的数据推送给客户端,而生成数据的阻塞操作(sleep)则在另一个线程里执行,互不干扰。
不用显式创建Thread的实时推送方案
Reactor本身提供了很多异步工具,完全不需要手动创建Thread,推荐这几种方式:
1. 使用Reactor调度器(Scheduler)转移阻塞任务
Reactor的Schedulers.boundedElastic()专门用来处理阻塞性任务,它会维护一个弹性线程池,自动管理线程生命周期。你可以用publishOn或者subscribeOn把生成数据的逻辑转移到这个线程池里:
@GetMapping public Flux<String> search() { return Flux.create(fluxSink -> { Random r = new Random(); for (int i = 0; i < 10; i++) { int n = r.nextInt(1000); System.out.println("Creating:" + n); fluxSink.next(String.valueOf(n)); try { Thread.sleep(2000); } catch (InterruptedException e) { e.printStackTrace(); } } fluxSink.complete(); }).publishOn(Schedulers.boundedElastic()); }
如果想更贴合Reactor的最佳实践,用Flux.generate(专门用于逐个生成元素的场景)配合subscribeOn会更规范:
@GetMapping public Flux<String> search() { return Flux.generate(sink -> { Random r = new Random(); int n = r.nextInt(1000); System.out.println("Creating:" + n); sink.next(String.valueOf(n)); try { Thread.sleep(2000); } catch (InterruptedException e) { e.printStackTrace(); } }) .take(10) // 只生成10个元素 .subscribeOn(Schedulers.boundedElastic()); }
2. 使用Flux.interval实现定时异步生成
如果你的需求是每隔固定时间生成一个数据,Flux.interval是最省心的选择——它本身就是异步的,会在专门的调度器线程里定时生成元素,完全不会阻塞EventLoop:
@GetMapping public Flux<String> search() { Random r = new Random(); return Flux.interval(Duration.ofSeconds(2)) // 每隔2秒生成一个元素 .take(10) // 生成10个后停止 .map(i -> { int n = r.nextInt(1000); System.out.println("Creating:" + n); return String.valueOf(n); }); }
核心总结
本质问题就是不要在EventLoop线程里执行阻塞操作,Reactor提供的调度器和异步操作符已经帮我们封装好了线程管理的逻辑,完全不需要手动创建Thread来处理异步推送。
内容的提问来源于stack exchange,提问作者Rodrigo De Oliveira Murta

