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

如何通过请求停止或暂停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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.04 16:18:03