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

Spring Boot中如何通过REST API从TaskExecutor返回动态响应?

Spring Boot动态进度REST API实现方案

针对你需要向客户端实时推送计数器进度的需求,有两种可行的实现方案,具体如下:

方案一:在/execute接口内直接返回实时进度(基于SSE)

这种方式利用**服务器发送事件(SSE)**实现单向实时推送,客户端发起一次请求后保持连接,服务器每隔10秒主动推送进度数据,直到任务结束。无需客户端反复请求,是更优雅的实时交互方式。

代码示例:

  1. 控制器实现:
@RestController
public class ProgressController {
    private final TaskExecutor taskExecutor;

    // 构造方法注入TaskExecutor
    public ProgressController(TaskExecutor taskExecutor) {
        this.taskExecutor = taskExecutor;
    }

    @GetMapping("/execute")
    public SseEmitter executeTask() {
        // 设置超时时间为5分钟+10秒,避免任务未完成就断开连接
        SseEmitter emitter = new SseEmitter(310000L);
        int totalSeconds = 300; // 5分钟=300秒
        int interval = 10; // 间隔10秒

        taskExecutor.execute(() -> {
            try {
                for (int current = 10; current <= totalSeconds; current += interval) {
                    Thread.sleep(10000); // 等待10秒
                    // 推送进度数据,包含当前已执行秒数和完成百分比
                    emitter.send(SseEmitter.event()
                            .data(Map.of("currentSecond", current, "progress", (current * 100) / totalSeconds)));
                }
                // 任务完成后发送结束标识
                emitter.send(SseEmitter.event()
                        .data(Map.of("status", "completed", "currentSecond", totalSeconds)));
                emitter.complete();
            } catch (Exception e) {
                emitter.completeWithError(e);
            }
        });

        return emitter;
    }
}
  1. 客户端侧(JavaScript示例):
const eventSource = new EventSource('/execute');
eventSource.onmessage = function(event) {
    const progress = JSON.parse(event.data);
    console.log(`当前进度:${progress.currentSecond}秒,完成${progress.progress}%`);
    if (progress.status === 'completed') {
        eventSource.close();
        console.log('任务已完成');
    }
};
eventSource.onerror = function(error) {
    console.error('连接出错', error);
    eventSource.close();
};

方案二:拆分为/execute和/get-progress两个接口

这种方式将任务启动和进度查询分离,客户端先调用/execute启动任务并获取唯一任务ID,之后定期调用/get-progress接口查询进度。

代码示例:

  1. 进度存储工具类:
@Component
public class TaskProgressStore {
    // 用并发HashMap存储任务ID与进度的映射
    private final ConcurrentHashMap<String, Integer> taskProgressMap = new ConcurrentHashMap<>();

    public void setProgress(String taskId, int progress) {
        taskProgressMap.put(taskId, progress);
    }

    public Integer getProgress(String taskId) {
        return taskProgressMap.get(taskId);
    }

    public void removeTask(String taskId) {
        taskProgressMap.remove(taskId);
    }
}
  1. 控制器实现:
@RestController
public class TaskController {
    private final TaskExecutor taskExecutor;
    private final TaskProgressStore progressStore;
    private final AtomicInteger taskIdGenerator = new AtomicInteger(0);

    // 构造方法注入依赖
    public TaskController(TaskExecutor taskExecutor, TaskProgressStore progressStore) {
        this.taskExecutor = taskExecutor;
        this.progressStore = progressStore;
    }

    @PostMapping("/execute")
    public ResponseEntity<String> startTask() {
        String taskId = "task_" + taskIdGenerator.incrementAndGet();
        int totalSeconds = 300;
        int interval = 10;

        taskExecutor.execute(() -> {
            try {
                for (int current = 10; current <= totalSeconds; current += interval) {
                    Thread.sleep(10000);
                    progressStore.setProgress(taskId, current);
                }
                progressStore.setProgress(taskId, totalSeconds);
                // 任务完成后延迟删除进度记录,避免客户端最后一次查询不到
                new Timer().schedule(() -> progressStore.removeTask(taskId), 5000);
            } catch (Exception e) {
                progressStore.removeTask(taskId);
            }
        });

        return ResponseEntity.ok(taskId);
    }

    @GetMapping("/get-progress")
    public ResponseEntity<Map<String, Object>> getProgress(@RequestParam String taskId) {
        Integer currentSecond = progressStore.getProgress(taskId);
        if (currentSecond == null) {
            return ResponseEntity.status(HttpStatus.NOT_FOUND)
                    .body(Map.of("message", "任务不存在或已完成"));
        }
        int totalSeconds = 300;
        return ResponseEntity.ok(Map.of(
                "currentSecond", currentSecond,
                "progress", (currentSecond * 100) / totalSeconds,
                "status", currentSecond == totalSeconds ? "completed" : "running"
        ));
    }
}
  1. 客户端侧(JavaScript示例):
// 启动任务
fetch('/execute', { method: 'POST' })
    .then(res => res.text())
    .then(taskId => {
        const interval = setInterval(() => {
            // 每隔10秒查询进度
            fetch(`/get-progress?taskId=${taskId}`)
                .then(res => res.json())
                .then(data => {
                    console.log(`当前进度:${data.currentSecond}秒,完成${data.progress}%`);
                    if (data.status === 'completed') {
                        clearInterval(interval);
                        console.log('任务已完成');
                    }
                })
                .catch(err => {
                    clearInterval(interval);
                    console.error('查询进度失败', err);
                });
        }, 10000);
    });

方案选择建议

  • 优先选方案一(SSE):实现简单,无需客户端轮询,实时性好,适合大多数现代浏览器环境。
  • 选方案二:兼容性更强(支持所有能发起HTTP请求的客户端),或者需要客户端主动控制查询频率的场景,但会增加服务器的请求量。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.05 13:42:37