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 进程
- 参数验证与身份认证。
- 创建数据库业务记录。
- 添加 BullMQ Job。
- 快速返回 Job ID。
- 提供业务状态查询接口。
Worker 进程
- 根据 ID 读取最新业务数据。
- 执行耗时操作。
- 更新数据库中的最终状态。
- 抛出错误,让 BullMQ 根据策略重试。
- 保证重复执行不会造成错误结果。
上线前检查
- Redis 配置
maxmemory-policy noeviction。 - 启用满足业务要求的 Redis 持久化。
- Producer 在 Redis 故障时快速返回
503,而不是无限占用 HTTP 连接。 - Worker 监听
error、failed和 stalled 相关事件。 - Job 数据不包含密码、令牌或大文件。
- 并发值经过真实压测。
- MySQL 中保存业务最终状态,Redis 只保存队列运行状态。
- Worker 处理函数具有幂等性。