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

如何用Reactor Core Java向浏览器推送长任务执行状态?

实现方案:用Reactor Core推送长任务的实时状态更新

Absolutely, you can pull this off with Reactor Core! Your goal—pushing a custom status object to the browser every time executeLongOperation() finishes—fits perfectly with reactive streams and Server-Sent Events (SSE), which your controller is already set up to use. Let's walk through the steps to adapt your code.

1. 定义自定义状态对象

First, create a POJO to represent the task status you want to send to the client. This should capture details like operation result, status (success/failure), and any relevant context:

public class TaskStatus {
    private String status; // 可选值:"SUCCESS", "FAILED", "IN_PROGRESS"
    private Object operationResult; // 存储executeLongOperation()的返回值
    private String message; // 可选:附加操作描述信息

    // 构造方法、getter/setter、toString()
    public TaskStatus(String status, Object operationResult, String message) {
        this.status = status;
        this.operationResult = operationResult;
        this.message = message;
    }

    // 省略getter和setter实现
}

2. 改造长任务为响应式Flux

Your original doLongTask is a synchronous loop—we need to convert this into a reactive Flux<TaskStatus> that emits a status update every time executeLongOperation() completes.

Key points to note:

  • Use Flux.fromIterable() to convert your List<Something> into a reactive stream.
  • For each Something in the list, process its conditional executeLongOperation() calls and emit a status for each execution.
  • Wrap blocking operations (like executeLongOperation()) in subscribeOn(Schedulers.boundedElastic()) to avoid blocking Reactor's main event loop—critical for long-running tasks.

Here's the adapted method:

public Flux<TaskStatus> doLongTaskReactive(List<Something> list) {
    return Flux.fromIterable(list)
        .flatMap(sm -> {
            Flux<TaskStatus> operations = Flux.empty();

            // 处理第一个条件下的长操作
            if (condition1) {
                operations = operations.concatWith(
                    Mono.fromCallable(() -> executeLongOperation())
                        .map(result -> new TaskStatus("SUCCESS", result, "完成条件1操作,目标ID:" + sm.getId()))
                        .onErrorResume(e -> Mono.just(new TaskStatus("FAILED", null, "条件1操作失败,目标ID:" + sm.getId() + ",错误:" + e.getMessage())))
                        .subscribeOn(Schedulers.boundedElastic()) // 把阻塞任务放到专用线程池
                );
            }

            // 处理第二个条件下的长操作
            if (condition2) {
                operations = operations.concatWith(
                    Mono.fromCallable(() -> executeLongOperation())
                        .map(result -> new TaskStatus("SUCCESS", result, "完成条件2操作,目标ID:" + sm.getId()))
                        .onErrorResume(e -> Mono.just(new TaskStatus("FAILED", null, "条件2操作失败,目标ID:" + sm.getId() + ",错误:" + e.getMessage())))
                        .subscribeOn(Schedulers.boundedElastic())
                );
            }

            return operations;
        });
}

为什么这么设计?

  • flatMap 允许我们处理每个Something对象,并为每个成功/失败的executeLongOperation()事件推送一条状态。
  • Mono.fromCallable() 将阻塞的executeLongOperation()调用包装为响应式类型。
  • onErrorResume 确保即使某个操作失败,整个流也不会中断,而是推送一条失败状态信息。
  • Schedulers.boundedElastic() 创建专用线程池处理阻塞任务,避免占用Reactor的非阻塞主线程。

3. 更新服务层方法

替换原来的占位方法,调用改造后的响应式长任务(注意根据实际业务逻辑获取你的List<Something>):

@Override
public Flux<TaskStatus> getTaskStatusReactor() {
    // 替换为你实际获取任务列表的逻辑
    List<Something> taskList = fetchYourBusinessTaskList();
    return doLongTaskReactive(taskList);
}

4. 调整控制器

你的控制器已经适配了SSE,只需要更新返回类型匹配新的Flux<TaskStatus>即可:

@GetMapping(path = "/taskStatusReactor", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
@ResponseBody
public Flux<TaskStatus> getTaskStatusReactor() {
    logger.debug("任务状态更新请求已初始化。");
    return simpleSearchService.getTaskStatusReactor();
}

完成!当客户端调用这个接口时,浏览器会持续接收TaskStatus对象(Spring Boot会自动序列化为JSON),每完成一次executeLongOperation()就推送一条更新。

关键注意事项

  • 背压处理:如果任务推送更新的速度快于浏览器消费速度,Reactor会默认缓存事件。对于超长时间任务,可以考虑添加onBackpressureBuffer()设置缓存上限,或者用onBackpressureDrop()丢弃旧更新(SSE通常能较好处理这类场景)。
  • 线程池管理:Schedulers.boundedElastic()是处理阻塞任务的最优选择,它会自动扩容线程数而不会压垮系统,避免用Schedulers.parallel()(它是为CPU密集型任务设计的)。
  • 流的完整性:当所有操作完成后,Flux会自动结束,SSE连接也会优雅关闭。

内容的提问来源于stack exchange,提问作者Carlos Sanchez

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.27 06:49:42