SmallRye Mutiny如何阻塞线程至满足条件并终止未完成线程
SmallRye Mutiny实现并发累计+提前终止方案
需求定义
- 并行启动3个异步任务,分别执行对应接口调用逻辑
- 实时累计各任务返回的计数值,一旦累计值超过99,立刻终止所有未完成任务,返回结果
"99+" - 若所有任务执行完成后累计值仍未超过99,返回整数类型的最终累计结果
原实现问题梳理
- 语法与变量可见性问题:Lambda引用主线程局部变量要求变量为final/effectively final,通过
Uni.createFrom().item()传参是值传递,子线程修改无法同步到主线程,导致完成标记asyncFlag永远无法更新 - 阻塞逻辑缺陷:主线程通过do-while空轮询判断状态,既会空耗CPU资源,也会因为标记无法更新陷入无限循环
- 代码逻辑错误:第三个任务的取消句柄被错误赋值给
cancellableThreads2变量,导致第三个任务无法被正常取消;共享累计变量totalAll无并发安全保护,多线程写入会出现竞态;子线程内调用await().indefinitely()会阻塞执行线程,违背响应式编程原则
正确实现代码
首先自定义一个标记异常用于触发提前终止逻辑,后续通过Mutiny的失败传播机制自动取消所有未完成任务,全程不需要手动维护取消句柄和轮询标记:
import io.smallrye.mutiny.Uni; import jakarta.ws.rs.core.Response; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.concurrent.atomic.AtomicInteger; import org.slf4j.Logger; import org.slf4j.LoggerFactory; // 自定义提前终止标记异常,无额外栈开销 class ValueOverflowException extends RuntimeException { @Override public Throwable fillInStackTrace() { // 重写栈采集方法,减少异常创建开销 return this; } } public Uni<Response> processRequest(RequestDTO request) { private static final Logger LOGGER = LoggerFactory.getLogger(当前类名.class); // 线程安全的原子累计器,初始值为0,替代原非线程安全的totalAll对象 AtomicInteger totalAccumulator = new AtomicInteger(0); DTOResponse resp = new DTOResponse(); // 注意:生产环境请使用容器托管的全局线程池,不要每次请求新建线程池 ExecutorService taskExecutor = Executors.newFixedThreadPool(3); // 定义三个异步任务,不提前触发订阅 Uni<Void> task1 = method1(request) .runSubscriptionOn(taskExecutor) .invoke(response -> { DTOResponse taskRes = response.readEntity(DTOResponse.class); int addVal = (Integer) taskRes.total; int currentTotal = totalAccumulator.addAndGet(addVal); LOGGER.info("任务1执行完成,新增值:{},当前累计值:{}", addVal, currentTotal); if (currentTotal > 99) { throw new ValueOverflowException(); } }) .replaceWithVoid(); Uni<Void> task2 = method2(request) .runSubscriptionOn(taskExecutor) .invoke(response -> { DTOResponse taskRes = response.readEntity(DTOResponse.class); int addVal = (Integer) taskRes.total; int currentTotal = totalAccumulator.addAndGet(addVal); LOGGER.info("任务2执行完成,新增值:{},当前累计值:{}", addVal, currentTotal); if (currentTotal > 99) { throw new ValueOverflowException(); } }) .replaceWithVoid(); Uni<Void> task3 = method3(request) .runSubscriptionOn(taskExecutor) .invoke(response -> { DTOResponse taskRes = response.readEntity(DTOResponse.class); int addVal = (Integer) taskRes.total; int currentTotal = totalAccumulator.addAndGet(addVal); LOGGER.info("任务3执行完成,新增值:{},当前累计值:{}", addVal, currentTotal); if (currentTotal > 99) { throw new ValueOverflowException(); } }) .replaceWithVoid(); // 组合三个任务,默认fail-fast模式:任意任务失败/抛出异常,自动取消其余未完成任务 return Uni.combine().all().unis(task1, task2, task3) .discardItems() // 所有任务正常完成,未触发阈值 .onItem().transform(unused -> { resp.total = totalAccumulator.get(); return Response.ok(resp).build(); }) // 捕获到提前终止异常,返回99+ .onFailure(ValueOverflowException.class).recoverWithItem(unused -> { resp.total = "99+"; return Response.ok(resp).build(); }) // 收尾逻辑:如果是自建线程池需要在这里关闭,托管线程池可跳过 .eventually(taskExecutor::shutdown); }
关键实现说明
- 用
AtomicInteger作为累计容器,天然支持多线程安全读写,同时解决了Lambda无法修改局部变量的语法限制 - 完全移除手动维护的完成标记、空轮询逻辑、手动取消句柄,依赖Mutiny原生的组合算子失败传播能力实现自动等待、自动取消,不会出现空耗CPU、句柄赋值错误的问题
- 全程使用响应式链式算子,移除子任务内的阻塞调用
await().indefinitely(),不会阻塞执行线程 - 自定义轻量级异常作为提前终止信号,关闭栈采集减少性能开销,正常流程和终止流程逻辑解耦
- 所有任务取消逻辑由Mutiny框架自动完成,当任意任务检测到累计值超阈值抛出异常时,框架会主动取消其余未完成任务的订阅,不需要手动调用cancel方法
内容的提问来源于stack exchange,提问作者RxGianYagami
相关产品推荐
相关产品推荐

