Background jobs

队列

用于延迟和高吞吐量处理的简单任务队列——QStash 适用于 Serverless,BullMQ 适用于自托管 Redis 部署。

概述

@nebutra/queue 是一个支持两种后端的提供商无关任务队列:

  • QStash(Upstash)——Serverless,基于 HTTP,无需管理基础设施。适合 Vercel 和边缘部署。
  • BullMQ——基于 Redis,自托管,高吞吐量。适合具有持久队列的专用服务器。

两种后端的应用代码完全相同——只有环境变量不同。

提供商自动检测

队列后端根据存在的环境变量自动选择:

优先级条件提供商
1QUEUE_PROVIDER 已设置按指定值(qstash / bullmq / memory
2存在 QSTASH_TOKENqstash
3存在 REDIS_URLbullmq
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 redis
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();

// 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?

目录