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

SpringBoot endpoint如何实现处理任务完成即实时返回结果

SpringBoot接口逐任务即时返回结果实现方案

要实现单任务完成就立刻推送结果给调用方,不需要等所有任务执行完再统一返回,核心是放弃默认的「一次请求对应一次完整响应」的模式,改用HTTP流式响应机制:任务完成一块就往响应输出流写一块数据、即时刷新缓冲区,保持连接直到所有任务执行完毕。你之前用@Async、ThreadPoolExecutor没生效,核心原因是要么等所有任务跑完才组装返回值,要么写数据后没刷新缓冲区导致数据卡在服务端响应缓存里。

以下是两种生产可用的原生实现,不需要引入额外中间件:


方案1:SseEmitter(服务端推送事件,最适配多结果分次推送场景)

SSE是SpringMVC原生支持的服务端推送组件,专门用于服务端主动分次向客户端推送数据,连接建立后服务端可以随时发送数据,客户端可以逐事件接收,不需要等待整个响应完成。

代码示例

import org.springframework.http.MediaType;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RestController;
import org.springframework.web.servlet.mvc.method.annotation.SseEmitter;
import java.io.IOException;
import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicInteger;

@RestController
public class TaskProcessController {
    // 生产环境请自定义ThreadPoolExecutor参数,不要用Executors默认的无界队列线程池,避免OOM
    private final ThreadPoolExecutor taskExecutor = new ThreadPoolExecutor(
            5,
            10,
            60L,
            TimeUnit.SECONDS,
            new LinkedBlockingQueue<>(100),
            new ThreadPoolExecutor.CallerRunsPolicy()
    );

    @GetMapping(value = "/task/process", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
    public SseEmitter processMultiTasks() {
        // 初始化SSE发射器,超时时间按业务最长任务时长设置,单位毫秒
        SseEmitter emitter = new SseEmitter(300000L);
        // 任务计数器,用于判断所有任务是否执行完成
        AtomicInteger taskCounter = new AtomicInteger(yourTaskList.size());

        taskExecutor.execute(() -> {
            try {
                // 替换为你实际的任务列表
                for (int i = 0; i < yourTaskList.size(); i++) {
                    int taskNo = i;
                    // 每个任务单独提交线程池执行
                    taskExecutor.execute(() -> {
                        try {
                            // 执行单个任务获取结果
                            Object taskResult = yourTaskList.get(taskNo).call();
                            // 任务完成立刻推送结果,客户端即时收到
                            emitter.send(SseEmitter.event()
                                    .id("task-" + taskNo)
                                    .name("taskSuccess")
                                    .data(taskResult));
                        } catch (Exception e) {
                            // 单个任务失败推送错误事件,不中断其他任务执行
                            try {
                                emitter.send(SseEmitter.event()
                                        .id("task-" + taskNo)
                                        .name("taskFail")
                                        .data("任务" + taskNo + "执行失败:" + e.getMessage()));
                            } catch (IOException ex) {
                                throw new RuntimeException(ex);
                            }
                        } finally {
                            // 所有任务执行完成后关闭连接
                            if (taskCounter.decrementAndGet() == 0) {
                                emitter.complete();
                            }
                        }
                    });
                }
            } catch (Exception e) {
                emitter.completeWithError(e);
            }
        });

        // 接口立刻返回发射器实例,保持连接打开等待后续数据推送
        return emitter;
    }
}

注意事项

  • 客户端需要按SSE协议逐事件读取数据:前端可以直接用原生EventSourceAPI监听事件,服务间调用使用支持流式响应的HTTP客户端(如WebClient、OkHttp流式模式)即可,先完成的任务结果会第一时间推送到客户端。
  • 如果服务部署在Nginx、网关后面,需要调整proxy_read_timeout等超时参数,避免长连接被中间件提前断开。
  • 跨域场景下和普通接口一样配置CORS规则即可。

方案2:StreamingResponseBody(适合无事件格式的纯流输出场景)

如果不需要SSE的事件结构化格式,只是想按任务完成顺序逐块输出响应内容,可以使用StreamingResponseBody直接操作响应输出流,核心是每次写完数据立刻刷新缓冲区。

代码示例

import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RestController;
import org.springframework.web.servlet.mvc.method.annotation.StreamingResponseBody;
import java.util.List;
import java.util.concurrent.*;

@RestController
public class StreamTaskController {
    private final ThreadPoolExecutor taskExecutor = new ThreadPoolExecutor(
            5,10,60L,TimeUnit.SECONDS,
            new LinkedBlockingQueue<>(100),
            new ThreadPoolExecutor.CallerRunsPolicy()
    );

    @GetMapping(value = "/task/stream", produces = "text/plain;charset=UTF-8")
    public StreamingResponseBody processTasksInStream() {
        return outputStream -> {
            // 提交所有任务,用Future接收执行结果
            List<Future<Object>> taskFutures = yourTaskList.stream()
                    .map(task -> taskExecutor.submit(task))
                    .toList();
            // 遍历获取结果,哪个任务先完成就先输出
            for (Future<Object> future : taskFutures) {
                Object result = future.get();
                outputStream.write(("任务执行结果:" + result + "\n").getBytes());
                // 关键操作:立刻刷新输出缓冲区,数据会立刻发送到客户端
                outputStream.flush();
            }
            // 所有任务写完后关闭流
            outputStream.close();
        };
    }
}

不推荐的替代方案

如果你的客户端完全不支持HTTP长连接、只能等待完整响应后再解析,可以改用「异步任务+回调」模式:接口收到请求后立刻返回全局唯一任务ID,每个任务执行完成后,服务端主动调用客户端提前预留的回调接口推送结果。这种方案需要客户端配合提供回调接口、维护任务状态,开发成本更高,没有特殊限制优先选上面两种流式方案。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.29 00:03:20