Background jobs
队列
用于延迟和高吞吐量处理的简单任务队列——QStash 适用于 Serverless,BullMQ 适用于自托管 Redis 部署。
概述
@nebutra/queue 是一个支持两种后端的提供商无关任务队列:
- QStash(Upstash)——Serverless,基于 HTTP,无需管理基础设施。适合 Vercel 和边缘部署。
- BullMQ——基于 Redis,自托管,高吞吐量。适合具有持久队列的专用服务器。
两种后端的应用代码完全相同——只有环境变量不同。
提供商自动检测
队列后端根据存在的环境变量自动选择:
| 优先级 | 条件 | 提供商 |
|---|---|---|
| 1 | QUEUE_PROVIDER 已设置 | 按指定值(qstash / bullmq / memory) |
| 2 | 存在 QSTASH_TOKEN | qstash |
| 3 | 存在 REDIS_URL | bullmq |
| 4 | 都不存在 | memory(仅限开发/测试——非持久化) |
配置步骤
QSTASH_TOKEN=""
QSTASH_CURRENT_SIGNING_KEY=""
QSTASH_NEXT_SIGNING_KEY=""
# API 网关的公开基础 URL
# QStash 将调用此 URL 投递任务
QSTASH_CALLBACK_BASE_URL="https://api.yourdomain.com"API 网关已在以下地址暴露 QStash Webhook:
POST /api/v1/queue/:queue/:type当 report.generate 任务就绪时,QStash 将 POST 到 https://api.yourdomain.com/api/v1/queue/report/generate。
QStash 需要可公开访问的 HTTPS 接口,无法投递到 localhost。本地开发时请使用内网穿透工具(如 ngrok)。
import { getQueue, createJob } from "@nebutra/queue";
const queue = await getQueue();
await queue.enqueue(
createJob("report", "generate", {
tenantId: "org_123",
reportType: "monthly",
})
);import { getQueue } from "@nebutra/queue";
const queue = await getQueue();
queue.registerHandler("report", "generate", async (job) => {
await generateReport(job.data);
});REDIS_URL="redis://localhost:6379"
# 可选:即使 QSTASH_TOKEN 也存在,也强制使用 BullMQ
QUEUE_PROVIDER="bullmq"# 本地开发
docker run -d -p 6379:6379 redis:7-alpine
# 或通过 pnpm dev(如果 Redis 已在 compose 文件中)
docker compose up redisimport { getQueue, createJob } from "@nebutra/queue";
const queue = await getQueue();
await queue.enqueue(
createJob("report", "generate", {
tenantId: "org_123",
reportType: "monthly",
})
);import { getQueue } from "@nebutra/queue";
const queue = await getQueue();
// BullMQ 启动 Worker 进程拉取任务
queue.registerHandler("report", "generate", async (job) => {
await generateReport(job.data);
});入队与处理模式
所有提供商的入队/处理模式完全一致:
import { getQueue, createJob } from "@nebutra/queue";
// --- 生产者(入队)---
const queue = await getQueue();
await queue.enqueue(
createJob(
"email", // 队列名称
"send", // 任务类型
{ // 任务数据(可序列化)
to: "[email protected]",
template: "invoice-paid",
data: { invoiceId: "inv_123", amount: 99.99 },
},
{ // 选项(可选)
tenantId: "org_123",
delay: 5000, // 投递前延迟(毫秒)
}
)
);
// --- 消费者(处理器)---
queue.registerHandler("email", "send", async (job) => {
await sendEmail(job.data.to, job.data.template, job.data.data);
});任务选项
createJob(queueName, jobType, data, {
tenantId: "org_123", // 将任务限定到某个租户
delay: 60_000, // 60 秒后投递
priority: 10, // 值越大越优先处理(仅 BullMQ)
jobId: "unique-id", // 按 ID 去重(防止重复入队)
})QStash vs BullMQ 对比
优点:
- 零基础设施——由 Upstash 完全托管
- 可在 Vercel、Netlify 等 Serverless 平台运行
- 内置指数退避重试
- 基于 HTTP——易于检查和调试
- 签名 Webhook——投递经过认证
缺点:
- 需要公开的 Webhook 接口(不能直接使用
localhost) - 不支持优先级队列
- 超高并发时单条消息成本较高
- 最大消息体:1 MB
适用于: Vercel 部署、中低流量队列、不想管理 Redis 的团队。
优点:
- 高吞吐量——每天数百万任务
- 支持优先级队列
- 适用于私有网络(无需公开接口)
- 丰富的任务生命周期事件(active、completed、failed、stalled)
- 支持延迟和可重复任务
缺点:
- 需要 Redis 实例
- Worker 必须是长期运行的进程(不支持 Serverless)
- 运维开销更大
适用于: 自托管部署、高流量处理、需要优先级或复杂调度的工作负载。
死信队列(DLQ)
耗尽所有重试次数的任务会被移入死信队列。
QStash 使用指数退避算法重试失败投递(默认最多 3 次)。所有重试耗尽后,任务被丢弃。可在 Upstash 控制台 中监控失败投递。
// 入队时配置重试次数
await queue.enqueue(
createJob("report", "generate", data, { maxRetries: 5 })
);BullMQ 将耗尽重试的任务移入 failed 队列。可以通过代码查看和重试:
import { getQueue } from "@nebutra/queue";
const queue = await getQueue();
// 列出失败任务
const failedJobs = await queue.getFailedJobs("report");
// 重试特定失败任务
await queue.retryJob("report", jobId);
// 清空失败队列(丢弃所有)
await queue.cleanFailed("report");Python 支持
@nebutra/queue 也提供 Python 客户端,用于微服务:
from _shared.queue import get_queue, create_job
queue = await get_queue()
# 入队
await queue.enqueue(create_job("report", "generate", {"tenant_id": "org_123"}))
# 处理
@queue.handler("report", "generate")
async def handle_report(job):
await generate_report(job.data["tenant_id"])相关文档
How is this guide?
在 GitHub 上编辑此页面
最后更新于