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
相关产品推荐
相关产品推荐

