原创

AI 内容流水线要不要上任务队列?同步执行、后台任务与 BullMQ 对比

针对 AI 内容流水线从小规模到多人协作的演进,比较同步请求、进程内后台任务和 BullMQ 队列在超时、重试、幂等、并发、可观测性及故障恢复上的取舍,并给出按任务规模选择架构的边界。

AI 内容工程 AI 内容流水线 BullMQ Node.js 任务队列 AI 总结
AI 内容工程专题 · 第 4/11 篇查看专题目录 →

一条典型的 AI 内容流水线可能包含:

抓取来源 → 清洗与去重 → AI 总结 → 自动分类 → 人工审核 → 发布

开发初期,最自然的写法是让管理后台发送请求,服务器依次完成所有步骤,最后把结果返回页面。这种同步执行简单直观,但随着来源增加、模型响应变慢或中途出现故障,一个请求会承担越来越多职责。

是不是应该马上引入 Redis 和 BullMQ?答案通常是否定的。任务队列解决的是可靠调度、并发控制和故障恢复问题,但也会增加部署组件和维护成本。更合理的做法,是根据任务规模从同步请求、进程内后台任务逐步演进到持久化队列。

一、方案一:一个请求执行完整流水线

app.post('/api/admin/run', async (req, res) => {
  try {
    const items = await fetchSources();
    const cleaned = deduplicate(items);
    const summarized = await summarize(cleaned);
    const classified = await classify(summarized);
    await publish(classified);
    res.json({ ok: true, count: classified.length });
  } catch (error) {
    res.status(500).json({ ok: false, error: error.message });
  }
});

当数据量较小、执行时间可控、只有管理员偶尔手动触发时,同步请求代码路径短、调试方便,也不需要额外组件。

它的问题是请求生命周期与整个任务绑定。浏览器、反向代理和应用服务器都可能有超时限制。HTTP 请求失败也不代表服务器工作一定停止,管理员再次点击后可能产生重复任务。如果所有步骤都在一个 try/catch 中,分类失败还可能使抓取结果无法落盘。因此,即使暂时使用同步模式,也应尽早分别保存每一步的结果。

二、方案二:请求立即返回,任务在进程内继续

const jobs = new Map();

app.post('/api/admin/run', (req, res) => {
  const jobId = crypto.randomUUID();
  jobs.set(jobId, { status: 'queued', step: 'waiting' });
  setImmediate(() => {
    runPipeline(jobId).catch(error => {
      jobs.set(jobId, { ...jobs.get(jobId), status: 'failed', error: error.message });
    });
  });
  res.status(202).json({ jobId });
});

后台拿到 jobId 后查询状态,就可以显示“正在抓取”“正在总结”或“处理失败”,不再长时间等待一个 HTTP 请求。这是同步执行到正式队列之间很实用的一步,适合单实例、低频任务。

但 Map 中的状态只存在内存里。应用重启、容器更新或进程崩溃后,任务记录会消失。多实例部署时,状态查询也可能落到另一个进程。因此它不是可靠队列。

三、方案三:使用 BullMQ 和 Redis

当任务需要持久化、自动重试或横向扩展时,可以使用 BullMQ。生产者只负责创建任务:

import { Queue } from 'bullmq';

const connection = { host: process.env.REDIS_HOST, port: Number(process.env.REDIS_PORT || 6379) };
const contentQueue = new Queue('content-pipeline', { connection });

const job = await contentQueue.add('summarize', { batchId }, {
  jobId: `summarize:${batchId}`,
  attempts: 3,
  backoff: { type: 'exponential', delay: 2000 }
});

Worker 负责执行:

import { Worker } from 'bullmq';

const worker = new Worker('content-pipeline', async job => {
  if (job.name === 'summarize') return summarizeBatch(job.data.batchId);
  throw new Error(`未知任务类型:${job.name}`);
}, { connection, concurrency: 2 });

示例中的重试次数、退避时间和并发数只是演示值。上线前需要结合模型限流、单次任务时长和服务器资源进行测试。

队列可以让任务独立保存,让 Worker 单独扩容并限制并发。失败任务可以自动重试并保留原因,管理员不必重新运行整个流水线。不过,拆分越细,状态流转越复杂;小型网站没有必要一开始就建立多个队列。

四、任务队列不等于 Worker Threads

Node.js 的 worker_threads 和 BullMQ Worker 名称相似,但解决的问题不同。worker_threads 适合 CPU 密集型 JavaScript 计算,对以网络 I/O 为主的网页抓取和模型 API 调用帮助有限。

内容流水线更需要的是任务持久化、重试与退避、并发限制、状态查询以及多进程消费,这些属于任务队列解决的范围。

五、无论是否上队列,都要先解决幂等

业务代码不能假设每个任务永远只运行一次。文章抓取可使用规范化后的原始 URL 作为唯一键;总结任务则可以根据文章 ID、正文哈希、提示词版本和模型名称构造任务键:

const jobKey = ['summary', article.id, article.contentHash, promptVersion, modelName].join(':');

这样,相同输入重复提交时不会产生两份相同结果;正文、提示词或模型变化时,又允许生成新版本。发布前还应检查稳定 ID 或 slug,避免重试生成两篇公开文章。

六、重试不是越多越好

适合重试的通常是网络中断、上游暂时不可用、明确可恢复的限流和短暂超时。API Key 无效、参数错误、输入超限和代码逻辑错误则不适合盲目重试。

function isRetryable(error) {
  return ['ETIMEDOUT', 'ECONNRESET', 'EAI_AGAIN'].includes(error.code) ||
    error.status === 429 || error.status >= 500;
}

真实错误格式应以供应商文档和实际响应为准。重试还应使用退避和抖动,避免多个任务同时失败后又集中请求上游服务。

七、保存中间结果,才能从失败处恢复

抓取完成后应立即保存原始数据。总结失败时,管理员可以对同一批数据重新执行总结,而不是重新访问全部来源。具体拆分方法可继续阅读Node.js 抓取与 AI 总结拆分实践。

每次模型调用建议记录模型配置 ID、模型名称、起止时间、状态、错误码、输入哈希、提示词版本和输出校验结果,但不要记录 API Key。结构化输出可参考GLM JSON 结构化输出;单个模型持续不可用时,再进入多模型重试、轮询与熔断流程。

八、如何决定现在是否需要 BullMQ

可以检查以下问题:

  • Web 请求是否经常接近代理超时?
  • 服务重启后,等待中的任务是否必须保留?
  • 是否需要多个 Worker 并发处理?
  • 是否需要自动重试、延迟任务或优先级?
  • 是否存在多个应用实例?
  • 是否需要查看历史任务和失败原因?

如果多数答案是否定的,进程内后台任务加持久化状态可能已经足够。如果任务不能丢失、需要跨进程消费或持续积压,队列的价值会明显增加。Redis 的备份、内存、连接、监控和升级也需要维护,不能忽略这部分成本。

九、推荐的渐进式演进路线

第一阶段使用同步执行,但拆分 fetch()、summarize()、classify() 和 publish()。模型调用可参考GLM-4.7-Flash Node.js 接入。

第二阶段让请求立即返回任务 ID,页面轮询状态,中间结果写入持久化存储。

第三阶段在可靠性需求增加后接入 BullMQ 与 Redis,优先迁移最慢、最容易失败的 AI 总结步骤。

第四阶段将 Web 服务和 Worker 分开部署,监控队列积压、失败量、处理时长与上游错误,并建立失败任务重放入口。更多内容可查看 AI 内容工程专题。

FAQ

小型网站需要一开始就使用 Redis 吗?

通常不需要。单实例、低频运行时,可以先采用后台任务和持久化状态。任务不能丢失或需要跨进程处理时,再引入 Redis 队列。

setImmediate() 能保证任务执行完成吗?

不能。进程退出、容器重启或程序崩溃时,内存中的任务可能丢失。

BullMQ 会自动避免文章重复发布吗?

不会替业务层完成所有去重。固定 job ID 可以降低重复入队,但抓取、总结和发布函数自身仍应设计为幂等。

AI 总结应该设置多少并发?

没有适用于所有模型和服务器的固定数值。应根据供应商限流、输入长度、任务耗时和本机资源逐步测试。

抓取、总结和分类应该放在同一个任务中吗?

早期可以放在一个任务中,但仍应保存每一步结果。需要单独重跑、使用不同并发或隔离失败时,再拆成不同任务。

模型调用失败后是否应该立刻切换模型?

应根据错误类型决定。临时网络故障可以有限重试;鉴权或参数错误不应通过轮询掩盖;持续不可用时才进入故障切换流程。

实测与内容说明

实测记录

  • 本文代码为架构示例,已按 Node.js ESM 语法进行静态审阅;未进行生产环境压测,不提供未经验证的吞吐量、耗时或可用性数据。
  • 重试次数、退避时间和并发数均为演示配置,上线前需按模型限流、任务时长与服务器资源单独验证。

参考资料

内容版本 1.0 · 审核:推荐智能手记 · 计划复审:2026-11-13