队列限流

并发控制回答“同时最多执行多少个任务”,限流回答“一段时间内最多开始多少个任务”。调用有配额的第三方 API、发送短信或邮件时,即使 Worker 并发不高,也可能在短时间内超过服务端限制。

Worker limiter

import { Worker } from "bullmq";

interface SmsJobData {
  phone: string;
  content: string;
}

const worker = new Worker<SmsJobData>(
  "sms",
  async (job): Promise<void> => {
    await sendSms(job.data.phone, job.data.content);
  },
  {
    connection: { host: "127.0.0.1", port: 6379 },
    concurrency: 20,
    limiter: {
      max: 10,
      duration: 1_000,
    },
  },
);

async function sendSms(phone: string, content: string): Promise<void> {
  console.log({ phone, content });
}

上例允许较高并发来隐藏网络等待,但整个队列在 1 秒内最多开始处理 10 个任务。被限流的任务仍处于 waiting 状态,不算失败。

Worker limiter 对同一队列是全局生效的:即使运行 10 个 Worker,也不是每个 Worker 每秒 10 个,而是所有 Worker 合计遵守该限制。

队列级全局限流

还可以直接在 Queue 上设置和管理全局限制:

import { Queue } from "bullmq";

const smsQueue = new Queue("sms", {
  connection: { host: "127.0.0.1", port: 6379 },
});

// 整个队列每秒最多处理 10 个任务
await smsQueue.setGlobalRateLimit(10, 1_000);

const config = await smsQueue.getGlobalRateLimit();
const ttl = await smsQueue.getRateLimitTtl();

console.log({ config, ttl });

// 确定不再需要时才移除
await smsQueue.removeGlobalRateLimit();

Worker 上配置的 limiter 不会覆盖队列全局限制;多个限制同时存在时,实际吞吐会受到更严格的一方约束。

根据 HTTP 429 手动限流

静态配置无法预知外部服务动态返回的 Retry-After。Worker 可以调用 worker.rateLimit(duration) 暂停队列取任务,并抛出专用的 RateLimitError,让当前任务回到等待状态,而不是当作普通失败消耗重试次数。

import { RateLimitError, Worker } from "bullmq";

interface ApiJobData {
  resourceId: string;
}

interface ApiResponse {
  status: number;
  retryAfterMs?: number;
}

const worker = new Worker<ApiJobData>(
  "external-api",
  async (job): Promise<void> => {
    const response = await callExternalApi(job.data.resourceId);

    if (response.status === 429) {
      const retryAfterMs = response.retryAfterMs ?? 5_000;
      await worker.rateLimit(retryAfterMs);
      throw new RateLimitError();
    }

    if (response.status >= 500) {
      throw new Error(`外部服务错误:${response.status}`);
    }
  },
  {
    connection: { host: "127.0.0.1", port: 6379 },
    limiter: { max: 100, duration: 1_000 },
  },
);

async function callExternalApi(resourceId: string): Promise<ApiResponse> {
  console.log(`调用资源 ${resourceId}`);
  return { status: 200 };
}

即使主要依赖手动限流,也必须在 Worker 选项中提供 limiter.max,BullMQ 会用它判断是否执行限流校验。

与 HTTP 接口请求限流的区别

两者保护的入口不同,通常需要同时存在:

类型 保护对象 典型位置
HTTP 请求限流 Web 服务、登录接口、防滥用 Express 中间件、API Gateway
BullMQ 队列限流 Worker、数据库、第三方服务配额 Worker / Queue 配置

HTTP 接口把任务加入队列成功,并不代表 Worker 可以无限速消费。反过来,Worker 限流也不能阻止恶意用户高速请求 API 并把 Redis 队列塞满。

并发和限流要一起设计

假设外部 API 最多每秒 20 次、单次平均耗时 2 秒:

先从保守值开始,用真实延迟、429 数量和队列等待时间调整。

实用技巧与最佳实践

本章小结

concurrency 控制同时工作量,limiter 控制单位时间吞吐。静态配额使用 Worker limiter 或 Queue 的全局限流;外部服务返回 429 时,使用 worker.rateLimit() 配合 RateLimitError。队列限流不能替代面向用户的 HTTP 请求限流。

官方文档