Express 与 Worker 项目实战

本章给出一个适合当前学习项目的最小目录方案:HTTP 服务接收报表请求,BullMQ Worker 在独立进程生成报表。

目标结构

nloop/
├─ app.ts
├─ server.ts
├─ queues/
│  ├─ connection.ts
│  └─ report-queue.ts
├─ workers/
│  └─ report-worker.ts
└─ routes/
   └─ report.ts

教程强调理解角色,因此暂不引入复杂框架或依赖注入容器。

1. Redis 连接配置

queues/connection.ts

import type { ConnectionOptions } from "bullmq";

export const producerConnection: ConnectionOptions = {
  host: process.env.REDIS_HOST ?? "127.0.0.1",
  port: Number(process.env.REDIS_PORT ?? 6379),
  enableOfflineQueue: false,
  maxRetriesPerRequest: 1,
};

export const workerConnection: ConnectionOptions = {
  host: process.env.REDIS_HOST ?? "127.0.0.1",
  port: Number(process.env.REDIS_PORT ?? 6379),
  maxRetriesPerRequest: null,
};

Producer 运行在 HTTP 请求中,Redis 不可用时应快速失败;Worker 是后台消费者,通常应持续等待 Redis 恢复。

BullMQ 默认使用 IORedis。官方当前也介绍了 Node-Redis 和 Bun adapter,但初学教程先使用默认连接方式,资料和示例最完整。

2. 定义任务数据和 Queue

queues/report-queue.ts

import { Queue } from "bullmq";
import { producerConnection } from "./connection";

export interface GenerateReportJobData {
  reportId: number;
  requestedBy: number;
}

export interface GenerateReportResult {
  fileKey: string;
}

export const reportQueue = new Queue<
  GenerateReportJobData,
  GenerateReportResult,
  "generate-report"
>("report-generation", {
  connection: producerConnection,
  defaultJobOptions: {
    attempts: 3,
    backoff: {
      type: "exponential",
      delay: 1_000,
    },
    removeOnComplete: {
      age: 24 * 60 * 60,
      count: 1_000,
    },
    removeOnFail: {
      age: 7 * 24 * 60 * 60,
      count: 5_000,
    },
  },
});

reportId 指向数据库记录。不要把整份报表数据或敏感用户资料复制到 Redis Job 中。

3. Express 路由添加任务

routes/report.ts

import { Router } from "express";
import { reportQueue } from "../queues/report-queue";

export const reportRouter = Router();

reportRouter.post("/", async (request, response, next) => {
  try {
    const reportId: number = Number(request.body.reportId);
    const requestedBy: number = Number(request.body.requestedBy);

    if (!Number.isInteger(reportId) || !Number.isInteger(requestedBy)) {
      response.status(400).json({ message: "参数必须是整数" });
      return;
    }

    const job = await reportQueue.add(
      "generate-report",
      { reportId, requestedBy },
      {
        jobId: `report-${reportId}`,
      },
    );

    response.status(202).json({
      message: "报表任务已进入队列",
      jobId: job.id,
    });
  } catch (error: unknown) {
    next(error);
  }
});

202 Accepted 表示请求已被接受处理,不代表报表已经生成。

固定 jobId 可以阻止同一个 ID 在队列中重复存在,但已自动删除的旧 Job 不再参与这个判断。关键业务仍需数据库状态或唯一约束。

4. Worker 执行任务

workers/report-worker.ts

import { Job, Worker } from "bullmq";
import { workerConnection } from "../queues/connection";
import type {
  GenerateReportJobData,
  GenerateReportResult,
} from "../queues/report-queue";

async function generateReport(
  job: Job<GenerateReportJobData, GenerateReportResult, "generate-report">,
): Promise<GenerateReportResult> {
  await job.updateProgress(10);

  const report = await loadReportFromDatabase(job.data.reportId);
  await job.updateProgress(40);

  const fileKey: string = await buildAndUploadReport(report);
  await job.updateProgress(100);

  return { fileKey };
}

const worker = new Worker<
  GenerateReportJobData,
  GenerateReportResult,
  "generate-report"
>("report-generation", generateReport, {
  connection: workerConnection,
  concurrency: 2,
});

worker.on("completed", (job) => {
  console.log("报表任务完成", { jobId: job.id });
});

worker.on("failed", (job, error: Error) => {
  console.error("报表任务失败", {
    jobId: job?.id,
    message: error.message,
  });
});

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

示例中的 loadReportFromDatabase()buildAndUploadReport() 代表业务函数,需要根据实际项目实现。

5. 单独启动 Worker

开发环境可以分别运行:

ts-node server.ts
ts-node workers/report-worker.ts

Worker 是长期运行进程。不要在每个 HTTP 请求中创建新 Worker。

可以在 package.json 中增加脚本,但这会修改项目配置,因此应在真正接入 BullMQ 时再做:

{
  "scripts": {
    "start": "nodemon --exec ts-node server.ts",
    "worker:report": "ts-node workers/report-worker.ts"
  }
}

6. 查询任务状态

reportRouter.get("/:jobId", async (request, response, next) => {
  try {
    const job = await reportQueue.getJob(request.params.jobId);

    if (job === undefined) {
      response.status(404).json({ message: "任务不存在或已被清理" });
      return;
    }

    response.json({
      id: job.id,
      state: await job.getState(),
      progress: job.progress,
      result: job.returnvalue,
      failedReason: job.failedReason,
    });
  } catch (error: unknown) {
    next(error);
  }
});

不要把 BullMQ Job 永久当作业务记录。任务可能根据清理策略被删除。报表的最终状态、文件地址和业务审计信息应保存到 MySQL。

7. 优雅关闭

Worker 进程:

let closing: boolean = false;

async function shutdown(signal: NodeJS.Signals): Promise<void> {
  if (closing) {
    return;
  }

  closing = true;
  console.log(`收到 ${signal},等待 Worker 关闭`);
  await worker.close();
}

process.once("SIGINT", () => {
  void shutdown("SIGINT");
});

process.once("SIGTERM", () => {
  void shutdown("SIGTERM");
});

HTTP 服务关闭时还要调用:

await reportQueue.close();

部署平台给予的终止宽限时间必须大于常见任务执行时间,否则未完成任务仍可能变成 stalled,并被其他 Worker 重新执行。

8. 推荐的职责边界

Express 进程

Worker 进程

上线前检查

官方资料