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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.28 12:48:18