如何在Java的Flux.subscribe中实现wait/notify机制?
Java反应式Flux的takeUntil控制与通知实现
首先明确:Reactor(Spring反应式框架)的takeUntil/takeUntilOther是基于信号驱动的操作,并非传统线程的wait/notify机制。强行用阻塞式wait会破坏反应式异步非阻塞特性,优先推荐原生反应式实现方式。
一、定时停止(固定等待时长后终止)
如果需求是等待指定时间后停止接收数据,直接用takeUntilOther搭配Mono.delay实现,无需手动处理线程阻塞:
import reactor.core.publisher.Flux; import java.time.Duration; // 获取数据Flux Flux<Record> rflux = query.sub(); // 等待5秒后自动停止接收,完成时执行通知逻辑 rflux.takeUntilOther(Mono.delay(Duration.ofSeconds(5))) .subscribe( record -> { /* 处理每条数据库记录 */ }, error -> error.printStackTrace(), // 错误处理 () -> { /* 完成通知:比如发送消息、更新状态等 */ System.out.println("数据接收已停止,通知完成"); } );
二、外部触发停止(自定义信号通知终止)
如果需要由外部事件(比如用户操作、其他线程信号)触发停止,用PublishSubject作为信号源,外部通过发送信号控制流的终止:
import reactor.core.publisher.Flux; import reactor.core.publisher.PublishSubject; // 创建停止信号源 PublishSubject<Void> stopTrigger = PublishSubject.create(); // 外部需要停止时,调用此方法发送信号:stopTrigger.onNext(null); // 绑定信号源到Flux Flux<Record> rflux = query.sub(); rflux.takeUntilOther(stopTrigger) .subscribe( record -> { /* 处理记录 */ }, Throwable::printStackTrace, () -> { /* 终止时的通知操作 */ System.out.println("收到停止信号,数据接收完成,已执行通知"); } );
三、兼容传统wait/notify(不推荐,仅作场景适配)
如果业务必须结合线程的wait/notify机制,可使用CountDownLatch实现,但注意这会阻塞反应式线程,违背异步设计原则:
import reactor.core.publisher.Flux; import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; CountDownLatch latch = new CountDownLatch(1); Flux<Record> rflux = query.sub(); rflux.takeUntil(record -> { try { // 等待latch被唤醒,最多等待10秒,超时自动停止 return latch.await(10, TimeUnit.SECONDS); } catch (InterruptedException e) { Thread.currentThread().interrupt(); return true; // 线程中断时停止接收 } }) .subscribe( record -> { /* 处理记录 */ }, error -> { error.printStackTrace(); latch.countDown(); // 错误时唤醒等待线程 }, () -> { latch.countDown(); // 正常完成时触发notify System.out.println("数据接收完成,已唤醒等待线程"); } ); // 其他线程中触发停止:latch.countDown();
额外优化:修正Session创建代码
你当前的Mono.just(getDataStore().getRxSession())是预加载Session,不符合反应式惰性加载原则,建议改为懒加载形式,避免提前创建资源:
return Flux.usingWhen( Mono.fromSupplier(() -> getDataStore().getRxSession()), // 懒加载Session session -> Flux.from(session.run(query).records()), RxSession::close );
内容的提问来源于stack exchange,提问作者Amit G
相关产品推荐
相关产品推荐

