Node.js/Express长耗时批量任务最佳实践及故障续跑咨询
基于Loopback/Express的异步任务+断点续传最佳实践
一、核心架构拆分
- 彻底拆分API请求处理与异步任务执行:API仅负责参数校验、生成任务记录存入数据库,随后立即返回
202 Accepted,附带任务ID供客户端查询进度。 - 任务执行逻辑单独抽离为Worker进程或任务队列,避免阻塞Express主进程的请求处理能力。
二、任务持久化与断点续传设计
- 任务表必须包含核心字段:
task_id(唯一标识)、status(pending/running/success/failed)、current_step(当前执行步骤)、progress(进度百分比)、error_msg(错误信息)、payload(任务原始参数JSON)、interim_data(中间执行结果)。 - 每完成一个步骤,立即更新数据库中任务的
current_step和progress,禁止批量延迟更新。若任务依赖中间结果,同步写入interim_data字段。 - 服务器启动/故障恢复时,Worker进程先查询数据库中
status为running或pending的任务,依据current_step跳过已完成步骤,直接从断点处重启执行。
三、异步任务执行的两种实现方案
方案1:内置Worker线程(适合轻量任务)
用Node.js worker_threads模块创建独立线程执行任务,主进程无需等待任务完成即可返回:
// Loopback控制器示例 const { Worker } = require('worker_threads'); const TaskModel = require('../models/task'); async function startTask(req, res) { const taskPayload = req.body; // 1. 初始化任务记录 const task = await TaskModel.create({ payload: taskPayload, status: 'pending', current_step: 0, progress: 0 }); // 2. 启动Worker线程 const worker = new Worker('./task-worker.js', { workerData: { taskId: task.id, payload: taskPayload } }); // 3. 监听Worker异常,更新任务状态 worker.on('error', (err) => { TaskModel.updateById(task.id, { status: 'failed', error_msg: err.message }); }); worker.on('exit', (code) => { if (code !== 0) { TaskModel.updateById(task.id, { status: 'failed', error_msg: `Worker退出码:${code}` }); } }); // 4. 立即返回响应 res.status(202).json({ taskId: task.id, message: '任务已启动' }); }
Worker线程内的任务执行逻辑:
// task-worker.js const { workerData } = require('worker_threads'); const TaskModel = require('../models/task'); async function runTask() { const { taskId, payload } = workerData; const steps = ['数据校验', '第三方API调用', '结果入库']; await TaskModel.updateById(taskId, { status: 'running' }); for (let i = 0; i < steps.length; i++) { try { // 执行当前步骤业务逻辑 await executeStep(steps[i], payload); // 更新任务进度 await TaskModel.updateById(taskId, { current_step: i + 1, progress: Math.round(((i + 1)/steps.length)*100) }); } catch (err) { await TaskModel.updateById(taskId, { status: 'failed', error_msg: err.message, current_step: i + 1 }); throw err; } } // 任务完成 await TaskModel.updateById(taskId, { status: 'success', progress: 100 }); } runTask().catch(err => console.error(err));
方案2:分布式任务队列(适合高并发/重量级任务)
使用BullMQ或Bee-Queue这类队列工具,将任务消息存入Redis,Worker进程监听队列执行任务,自带重试、持久化、分布式部署支持:
// Loopback控制器示例 const Queue = require('bullmq').Queue; const taskQueue = new Queue('long-tasks', { connection: { host: 'localhost', port: 6379 } }); const TaskModel = require('../models/task'); async function startTask(req, res) { const taskPayload = req.body; // 1. 创建任务记录 const task = await TaskModel.create({ payload: taskPayload, status: 'pending', current_step: 0, progress: 0 }); // 2. 添加任务到队列 await taskQueue.add('execute-task', { taskId: task.id, payload: taskPayload }); // 3. 返回响应 res.status(202).json({ taskId: task.id, message: '任务已启动' }); }
Worker进程监听队列:
const { Worker } = require('bullmq'); const TaskModel = require('../models/task'); const worker = new Worker('long-tasks', async (job) => { const { taskId, payload } = job.data; const steps = ['数据校验', '第三方API调用', '结果入库']; await TaskModel.updateById(taskId, { status: 'running' }); for (let i = 0; i < steps.length; i++) { await executeStep(steps[i], payload); await TaskModel.updateById(taskId, { current_step: i + 1, progress: Math.round(((i + 1)/steps.length)*100) }); // 更新队列任务进度(可选,用于监控) job.updateProgress(Math.round(((i + 1)/steps.length)*100)); } await TaskModel.updateById(taskId, { status: 'success', progress: 100 }); }, { connection: { host: 'localhost', port: 6379 } });
四、进度查询与Long Polling实现
- 在Loopback中创建
getTaskProgress接口,根据taskId查询数据库返回任务状态:
async function getTaskProgress(req, res) { const { taskId } = req.params; const task = await TaskModel.findById(taskId); if (!task) { return res.status(404).json({ message: '任务不存在' }); } res.json({ taskId: task.id, status: task.status, progress: task.progress, currentStep: task.current_step, errorMsg: task.error_msg }); }
- 客户端实现Long Polling:若任务未完成,延迟3-5秒后重新发起请求;任务完成或失败则停止轮询。
五、故障恢复关键细节
- 任务幂等性:每个步骤需保证重复执行无副作用,比如使用唯一业务标识、数据库事务、幂等键。
- 定时巡检:Worker启动时或每分钟扫描数据库,标记长时间未更新的
running任务为failed,或自动重启执行。 - 原子更新:任务步骤执行与状态更新需用数据库事务包裹,避免执行成功但状态未更新的不一致情况。
六、Loopback专属优化
- 用Remote Hooks或Interceptors统一处理任务参数校验,减少重复代码。
- 封装任务模型的更新逻辑为Repository方法(如
updateTaskProgress、markTaskFailed),提升代码复用性。 - 利用
@Model装饰器定义任务模型字段与验证规则,保证数据一致性。
内容的提问来源于stack exchange,提问作者Demiro-FE-Architect
相关产品推荐
相关产品推荐

