如何阻塞等待热Flux发射的下一条数据并赋值?
解决方案
要实现阻塞等待热Flux的下一个发射数据,最简洁高效的方式是使用Flux.next()方法配合block():
// 假设hotFlux是你的热Flux<Integer>实例 Integer nextValue = hotFlux.next().block();
原理说明
blockFirst()之所以失效,是因为热Flux不会重放已发射的元素,它只会向订阅者推送订阅之后产生的新元素。blockFirst()会等待Flux的第一个元素,但如果热Flux已经发射过元素,这个方法会无限阻塞,因为不会再有“第一个”元素出现。next()方法会返回一个Mono<Integer>,它专门捕获订阅之后热Flux发射的第一个新元素,拿到元素后自动取消订阅,避免资源浪费。后续调用block()就能阻塞当前线程,直到这个元素到来。
优化建议
为了避免因热Flux长时间不发射数据导致的无限阻塞,建议添加超时时间:
import java.time.Duration; Integer nextValue = hotFlux.next().block(Duration.ofSeconds(10)); // 如果超时,nextValue会为null,可根据业务逻辑处理
内容的提问来源于stack exchange,提问作者rohan
相关产品推荐
相关产品推荐

