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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.24 14:25:26