13 - 隔离处理器、任务取消与超时

本章解决三个容易混淆的问题:什么时候普通 Worker 就足够,什么时候要使用 sandboxed processor(隔离处理器),以及“取消”和“超时”究竟能做到什么。

1. 先判断任务属于哪一类

Node.js 的 JavaScript 默认运行在事件循环线程上。如果处理器长时间占用 CPU,BullMQ 就可能没有机会及时续订任务锁,最终把任务判断为 stalled

实用判断方法:记录单次任务耗时,同时观察事件循环延迟和 stalled 事件。不要仅凭“任务耗时长”判断它是 CPU 密集型,因为一个耗时 30 秒的 HTTP 请求仍然可能主要是在等待 I/O。

2. 普通 Worker 的并发适合 I/O 任务

import { Job, Worker } from 'bullmq';

interface DownloadJobData {
  url: string;
}

const worker = new Worker<DownloadJobData>(
  'downloads',
  async (job: Job<DownloadJobData>): Promise<number> => {
    const response = await fetch(job.data.url);

    if (!response.ok) {
      throw new Error(`下载失败:HTTP ${response.status}`);
    }

    const body = await response.arrayBuffer();
    return body.byteLength;
  },
  {
    connection: { host: '127.0.0.1', port: 6379 },
    concurrency: 10,
  },
);

worker.on('error', (error: Error) => {
  console.error('Worker 错误', error);
});

concurrency 是单个 Worker 内的并发,不等于 CPU 并行。它利用事件循环在一个任务等待 I/O 时推进其他任务。应该从较小值开始,用生产负载压测后调整。

3. CPU 密集任务使用 sandboxed processor

隔离处理器把业务处理函数放到另一个 Node.js 运行环境中,使 BullMQ 的任务锁维护与 CPU 密集代码隔离。默认隔离方式是子进程;也可以启用 Worker Threads。

处理器应放在独立文件中:

// processors/calculate-report.ts
import type { SandboxedJob } from 'bullmq';

interface CalculateReportData {
  values: number[];
}

async function calculateReport(
  job: SandboxedJob<CalculateReportData>,
): Promise<number> {
  return job.data.values.reduce((sum: number, value: number) => sum + value, 0);
}

export = calculateReport;

Worker 接收的是处理器文件,而不是当前进程中的函数:

// report-worker.ts
import path from 'node:path';
import { pathToFileURL } from 'node:url';
import { Worker } from 'bullmq';

const processorPath = path.join(
  __dirname,
  'processors',
  'calculate-report.js',
);

const worker = new Worker(
  'reports',
  pathToFileURL(processorPath),
  {
    connection: { host: '127.0.0.1', port: 6379 },
    concurrency: 2,
  },
);

官方文档特别建议 Windows 使用 URL 形式的处理器路径。这里使用 TypeScript 的 export =,编译为 CommonJS 后会得到 BullMQ 官方示例要求的 module.exports = processor。还要注意 Worker 引用的是编译后的 .js 文件:隔离进程不能理所当然地继承主进程的 ts-node 注册状态。学习阶段可以先编译再启动 Worker,避免把模块加载问题误判为 BullMQ 问题。

如果希望使用 Worker Threads:

const worker = new Worker(
  'reports',
  pathToFileURL(processorPath),
  {
    connection: { host: '127.0.0.1', port: 6379 },
    concurrency: 2,
    useWorkerThreads: true,
  },
);

Worker Threads 通常比子进程少一些资源开销,但每个线程仍需要独立的 V8 运行时,并不是“几乎免费”。CPU 密集任务的进程或线程数量一般从 CPU 核心数附近开始压测,而不是设置成几十或几百。

4. 使用 AbortSignal 协作式取消

BullMQ 的处理函数可以通过第三个参数接收 AbortSignal。取消是协作式的:BullMQ 发出信号后,业务代码必须停止请求、终止循环并释放资源。

import { Job, UnrecoverableError, Worker } from 'bullmq';

interface ExportJobData {
  sourceUrl: string;
}

const worker = new Worker<ExportJobData>(
  'exports',
  async (
    job: Job<ExportJobData>,
    _token?: string,
    signal?: AbortSignal,
  ): Promise<string> => {
    try {
      const response = await fetch(job.data.sourceUrl, { signal });

      if (!response.ok) {
        throw new Error(`上游返回 HTTP ${response.status}`);
      }

      return await response.text();
    } catch (error: unknown) {
      if (signal?.aborted) {
        throw new UnrecoverableError(`任务被取消:${String(signal.reason)}`);
      }

      throw error;
    }
  },
  { connection: { host: '127.0.0.1', port: 6379 } },
);

worker.cancelJob('job-id-123', '用户取消了导出');

官方推荐事件驱动或把同一个 signal 传给原生支持取消的 API,而不是频繁轮询。对不支持 AbortSignal 的库,需要在 abort 监听器里调用它自己的 cancel()close() 等方法。

取消后抛出普通 Error,任务在还有 attempts 时会重试;抛出 UnrecoverableError 则不会重试。应根据业务语义选择,不能把所有取消都永久失败。

5. 为普通异步任务实现超时

BullMQ 官方的 timeout pattern 是在处理函数里使用 AbortController 和定时器。BullMQ 没有一个能够强行安全中断任意 JavaScript 的通用任务超时开关。

import { Job, UnrecoverableError, Worker } from 'bullmq';

interface FetchJobData {
  timeoutMs: number;
  url: string;
}

const worker = new Worker<FetchJobData>(
  'remote-fetch',
  async (job: Job<FetchJobData>): Promise<string> => {
    const controller = new AbortController();
    const timer = setTimeout(() => controller.abort('timeout'), job.data.timeoutMs);

    try {
      const response = await fetch(job.data.url, { signal: controller.signal });
      return await response.text();
    } catch (error: unknown) {
      if (controller.signal.aborted) {
        throw new UnrecoverableError(`任务超过 ${job.data.timeoutMs}ms`);
      }

      throw error;
    } finally {
      clearTimeout(timer);
    }
  },
  { connection: { host: '127.0.0.1', port: 6379 } },
);

如果超时属于临时故障,希望稍后重试,就抛出普通 Error;只有确定重试没有意义时才使用 UnrecoverableError

6. 隔离处理器的“硬超时”风险

隔离进程可以在超时后结束进程,但强制终止可能发生在写文件、提交事务或调用外部系统的中途,造成损坏或部分成功。官方 timeout-for-sandboxed-processors pattern 建议使用两阶段策略:

  1. 软超时:先尝试取消工作并清理资源。
  2. 硬超时:清理仍未结束时才终止隔离进程。

还有一个关键限制:如果处理器陷入阻塞事件循环的无限循环,写在同一个处理器里的 setTimeout 也无法执行。因此真正的硬超时监控必须位于被监控代码之外。

7. 最佳实践清单

官方资料