如何通过请求停止或暂停Fastify服务中Node.js运行的耗时函数与循环
实现方案
核心思路
你的原始代码是一次性生成所有间隔200ms的定时任务塞进事件循环,要实现停止/暂停能力,核心是两点:
- 把所有定时任务的ID存储起来,停止/暂停时可以直接清除未执行的定时任务
- 在Fastify实例上挂载全局的流处理状态,标记当前运行状态、处理进度、待处理数据等,供控制接口读取修改
具体实现步骤
1. 挂载全局流处理状态
在Fastify初始化逻辑中添加如下代码,把状态挂载到Fastify实例上,全局可访问:
// 初始化流处理全局状态 fastify.decorate('streamProcessing', { isRunning: false, // 是否有正在运行的处理任务 isPaused: false, // 是否处于暂停状态 timeoutIds: [], // 存储所有未执行的setTimeout ID currentIndex: 0, // 当前已处理到的流索引 pendingStreams: [] // 待处理的完整流列表,暂停恢复时使用 })
2. 修改原有的renderStreams逻辑
调整原逻辑,把定时任务ID存入状态,每次执行任务前先检查运行状态:
const renderStreams = (fastify, streams = []) => { const { redis, streamProcessing } = fastify; const channel = "streams"; // 清除之前残留的定时任务,避免冲突 streamProcessing.timeoutIds.forEach(id => clearTimeout(id)); // 重置状态 Object.assign(streamProcessing, { isRunning: true, isPaused: false, currentIndex: 0, pendingStreams: streams, timeoutIds: [] }) for (let i = 0; i < streams.length; i++) { const timeoutId = setTimeout(async () => { // 任务执行前先检查状态,已停止则直接跳过 if (!streamProcessing.isRunning) return; // 已暂停则记录当前进度后跳过 if (streamProcessing.isPaused) { streamProcessing.currentIndex = i; return; } const stream = streams[i]; await renderStream(redis, channel, stream); // 处理完成更新进度 streamProcessing.currentIndex = i + 1; }, i * 200) // 存储定时任务ID streamProcessing.timeoutIds.push(timeoutId); } }
3. 实现控制接口
完全停止接口
fastify.get('/api/streams/stop', async (req, reply) => { const { streamProcessing } = fastify; // 清除所有未执行的定时任务 streamProcessing.timeoutIds.forEach(id => clearTimeout(id)); // 重置所有状态 Object.assign(streamProcessing, { isRunning: false, isPaused: false, currentIndex: 0, pendingStreams: [], timeoutIds: [] }) return { success: true, message: '流处理已完全停止' } })
可选:暂停/恢复接口
如果需要支持暂停后继续处理,新增以下两个接口:
// 暂停处理 fastify.get('/api/streams/pause', async (req, reply) => { const { streamProcessing } = fastify; if (!streamProcessing.isRunning || streamProcessing.isPaused) { return { success: false, message: '无正在运行的处理任务,或已处于暂停状态' } } // 清除未执行的定时任务 streamProcessing.timeoutIds.forEach(id => clearTimeout(id)); streamProcessing.isPaused = true; return { success: true, message: '流处理已暂停', currentProcessedCount: streamProcessing.currentIndex } }) // 恢复处理 fastify.get('/api/streams/resume', async (req, reply) => { const { redis, streamProcessing } = fastify; if (!streamProcessing.isPaused || streamProcessing.pendingStreams.length === 0) { return { success: false, message: '无暂停中的任务,或无待处理流数据' } } const channel = "streams"; // 截取未处理的流列表 const remainingStreams = streamProcessing.pendingStreams.slice(streamProcessing.currentIndex); streamProcessing.isPaused = false; streamProcessing.timeoutIds = []; // 重新生成剩余流的定时任务 for (let i = 0; i < remainingStreams.length; i++) { const timeoutId = setTimeout(async () => { if (!streamProcessing.isRunning || streamProcessing.isPaused) return; const stream = remainingStreams[i]; await renderStream(redis, channel, stream); streamProcessing.currentIndex += 1; }, i * 200) streamProcessing.timeoutIds.push(timeoutId); } return { success: true, message: '流处理已恢复', remainingCount: remainingStreams.length } })
注意事项
- 当前实现不会中断已经开始执行的
renderStream任务,如果需要中断正在处理的任务,需要在renderStream内部每一步执行前都检查streamProcessing.isRunning状态,同时如果用到了ffmpeg子进程,停止时要主动kill对应的子进程,避免资源残留 - 该方案是单实例部署场景的实现,如果是多实例部署,需要把状态存储到Redis等共享存储中,同时做好任务唯一标记,避免多实例重复处理
- 可以根据需要加分布式锁,避免同时触发多个控制请求导致状态错乱
内容的提问来源于stack exchange,提问作者kinx
相关产品推荐
相关产品推荐

