队列限流
并发控制回答“同时最多执行多少个任务”,限流回答“一段时间内最多开始多少个任务”。调用有配额的第三方 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 秒:
limiter: { max: 20, duration: 1000 }控制吞吐。concurrency决定多少请求可同时等待响应。- 单纯把
concurrency设为 20,不等于每秒只调用 20 次。 - 单纯限流却把并发设为 1,可能无法充分利用允许的配额。
先从保守值开始,用真实延迟、429 数量和队列等待时间调整。
实用技巧与最佳实践
- 配额有突发和平均值两种口径时,以对方文档为准并留安全余量。
- 优先读取并遵守
Retry-After,不要把所有 429 固定等待 1 秒。 - 监控 waiting 数量、任务等待时长、429 次数和限流 TTL。
- 不要频繁调用
removeGlobalRateLimit()绕开保护;这会重置限流计数。 - 免费版的 limiter 是队列级全局限制,不要误以为可按
customerId分组;旧版 group key 支持已移除。 - 多租户需要独立配额时,可拆分队列,或评估 BullMQ Pro 的 Groups 能力。
- 限流任务回到 waiting 后仍可能重复调用外部系统,因此业务操作仍需幂等。
本章小结
concurrency 控制同时工作量,limiter 控制单位时间吞吐。静态配额使用 Worker limiter 或 Queue 的全局限流;外部服务返回 429 时,使用 worker.rateLimit() 配合 RateLimitError。队列限流不能替代面向用户的 HTTP 请求限流。