api-queue-bullmq
Job queues, background processing, and task scheduling with BullMQ v5
Install
npx skills add https://github.com/agents-inc/skills/tree/main/dist/plugins/api-queue-bullmq/skills/api-queue-bullmq
claude plugin marketplace add https://llmmart.ai/marketplace.json && claude plugin install agents-inc-skills@llmmart
git clone https://github.com/agents-inc/skills.git
The skills CLI installs just this skill, for any of its supported agents. Claude Code installs the whole agents-inc/skills collection as a plugin from our marketplace. Git is the plain clone.
Skill manifest
BullMQ Patterns
Quick Guide: Use BullMQ (v5.x) for background job processing, task scheduling, and workflow orchestration on top of Redis. Core classes:
Queue(adds jobs),Worker(processes jobs),QueueEvents(global event listener),FlowProducer(parent-child job trees). Always pass aconnectionobject to every constructor (required in v5). SetmaxRetriesPerRequest: nullon ioredis connections for Workers. UseupsertJobSchedulerfor repeatable/cron jobs (replaces deprecated repeatable API). QueueScheduler was removed in v4 -- its responsibilities are now handled by Workers automatically.
<critical_requirements>
CRITICAL: Before Using This Skill
All code must follow project conventions in CLAUDE.md (kebab-case, named exports, import ordering,
import type, named constants)
(You MUST pass a connection object to every Queue, Worker, QueueEvents, and FlowProducer constructor -- BullMQ v5 throws if connection is missing)
(You MUST set maxRetriesPerRequest: null on ioredis connections used by Workers -- BullMQ requires infinite retries and throws without this setting)
(You MUST call await worker.close() on SIGTERM/SIGINT for graceful shutdown -- without it, in-progress jobs become stalled)
(You MUST use upsertJobScheduler for repeatable/cron jobs -- the old repeat option on queue.add is deprecated since v5.16.0)
</critical_requirements>
Examples
- Core Patterns -- Queue setup, Worker processing, job options, connection factory, graceful shutdown, typed jobs
- Advanced Patterns -- FlowProducer, rate limiting, job scheduling, QueueEvents, concurrency, sandboxed processors
Additional resources:
- reference.md -- Decision frameworks, job option reference, anti-patterns, production checklist
Auto-detection: BullMQ, bullmq, Queue, Worker, QueueEvents, FlowProducer, job queue, background job, worker process, job scheduler, upsertJobScheduler, rate limiter, job priority, job delay, sandboxed processor, repeatable job, cron job, flow producer, parent child jobs
When to use:
- Background processing (email sending, image processing, PDF generation)
- Scheduled/cron jobs (nightly reports, periodic cleanup)
- Workflow orchestration with parent-child job dependencies (FlowProducer)
- Rate-limited API consumption (throttling outbound requests)
- Priority-based job processing (urgent jobs before bulk operations)
- Distributing CPU-intensive work across multiple workers or machines
Key patterns covered:
- Queue and Worker setup with typed job data and return values
- Connection factory with
maxRetriesPerRequest: nullfor Workers - Job options: delay, priority, attempts, backoff, removeOnComplete/Fail
- Graceful shutdown with
worker.close()on process signals - FlowProducer for parent-child job trees with dependency tracking
- Job Schedulers for repeatable/cron jobs (
upsertJobScheduler) - Rate limiting (global limiter and manual
Worker.RateLimitError) - Concurrency control (local per-worker and global)
- QueueEvents for global event monitoring across all workers
- Sandboxed processors for CPU-intensive work
When NOT to use:
- Simple in-process timers or
setTimeout(no persistence needed) - Real-time pub/sub messaging without persistence (use your pub/sub solution)
- Data that must be processed synchronously within a request-response cycle
- Queues that don't need persistence, retries, or scheduling
<decision_framework>
Decision Framework
When to Use BullMQ
Do you need background job processing?
|-- NO -> Don't use BullMQ
+-- YES -> Do you need persistence, retries, or scheduling?
|-- NO -> Simple in-process queue or setTimeout may suffice
+-- YES -> Do you need parent-child job dependencies?
|-- YES -> BullMQ with FlowProducer
+-- NO -> Do you need rate limiting or priority?
|-- YES -> BullMQ with limiter/priority options
+-- NO -> BullMQ with basic Queue + Worker
Which Job Pattern?
What kind of job scheduling do you need?
|-- One-time delayed job -> queue.add() with delay option
|-- Recurring on fixed interval -> upsertJobScheduler with every
|-- Recurring on cron schedule -> upsertJobScheduler with pattern
|-- Job that depends on other jobs -> FlowProducer with children
|-- Bulk of independent jobs -> queue.addBulk([...])
Concurrency Strategy
Is the processor CPU-intensive?
|-- YES -> Use sandboxed processor (file path or useWorkerThreads)
+-- NO -> Is it I/O-bound (network calls, DB queries)?
|-- YES -> Set concurrency option (e.g., 10-50)
+-- NO -> Default concurrency (1) is fine
</decision_framework>
<red_flags>
RED FLAGS
High Priority Issues:
- Missing
connectionon Queue/Worker/QueueEvents constructor -- BullMQ v5 throws at startup without it - Missing
maxRetriesPerRequest: nullon Worker connections -- BullMQ throws immediately - No graceful shutdown handler -- in-progress jobs become stalled on process exit and are re-processed by other Workers
- Using deprecated
repeatoption onqueue.add()instead ofupsertJobScheduler-- deprecated since v5.16.0 - Using integer job IDs -- BullMQ v5 throws; IDs must be strings
Medium Priority Issues:
- No
removeOnComplete/removeOnFailconfigured -- completed/failed jobs accumulate in Redis indefinitely - Sharing a single ioredis connection between Worker and QueueEvents -- both use blocking commands and will interfere
- CPU-intensive processor without sandboxed mode -- blocks event loop, causes stalled jobs, prevents lock renewal
- Missing error event handler on Worker --
worker.on("error", ...)prevents unhandled errors from crashing the process
Common Mistakes:
- Assuming
worker.close()has a timeout -- it waits indefinitely for processors to finish; wrap with your own timeout - Using QueueScheduler class -- removed in BullMQ v4; its responsibilities are handled by Workers automatically
- Expecting
limiterto be per-Worker -- the rate limit is global across all Workers on the same queue - Not making processors idempotent -- BullMQ guarantees at-least-once delivery; a stalled job may be processed twice
Gotchas & Edge Cases:
upsertJobSchedulereveryintervals align to the clock (0s, 2s, 4s), not to when you called the methodworker.close()does not cancel running processors -- it waits for them to finish naturally- Job priority has a performance cost -- BullMQ uses a different data structure for priority queues; skip if not needed
QueueEventsuses Redis Streams internally -- ensure your Redis instance has sufficient memory for stream data- FlowProducer adds the entire tree atomically -- if any child fails validation, none are added
job.getChildrenValues()returns an object keyed by"queueName:jobId"-- not an array- Redis must have
maxmemory-policyset tonoeviction-- BullMQ relies on keys not being evicted
</red_flags>
<critical_reminders>
CRITICAL REMINDERS
All code must follow project conventions in CLAUDE.md (kebab-case, named exports, import ordering,
import type, named constants)
(You MUST pass a connection object to every Queue, Worker, QueueEvents, and FlowProducer constructor -- BullMQ v5 throws if connection is missing)
(You MUST set maxRetriesPerRequest: null on ioredis connections used by Workers -- BullMQ requires infinite retries and throws without this setting)
(You MUST call await worker.close() on SIGTERM/SIGINT for graceful shutdown -- without it, in-progress jobs become stalled)
(You MUST use upsertJobScheduler for repeatable/cron jobs -- the old repeat option on queue.add is deprecated since v5.16.0)
Failure to follow these rules will cause startup crashes, stalled jobs, and unreliable job processing.
</critical_reminders>
Files (skills)
-
examples
-
advanced.md 9.7 KB
# BullMQ Advanced Patterns > Related: [core.md](core.md) -- Queue setup, Worker processing, job options, connection factory, graceful shutdown --- ## FlowProducer (Parent-Child Job Trees) FlowProducer creates DAG-like job trees where a parent waits for all children to complete before being processed itself. The entire tree is added atomically. ```typescript // ✅ Good -- Video processing pipeline with FlowProducer import { FlowProducer, Worker, type Job } from "bullmq"; const flowProducer = new FlowProducer({ connection: createWorkerConnection() }); // Add a flow -- parent waits for all children const flow = await flowProducer.add({ name: "publish-video", queueName: "videos", data: { videoId: "abc-123" }, children: [ { name: "extract-audio", queueName: "media-processing", data: { videoId: "abc-123" }, }, { name: "generate-thumbnail", queueName: "media-processing", data: { videoId: "abc-123", size: "720p" }, }, { name: "transcode-hd", queueName: "media-processing", data: { videoId: "abc-123", format: "mp4" }, }, ], }); // flow.job = parent job reference // flow.children = child job references ``` **Why good:** Atomic addition (all jobs added or none), parent automatically waits for all children, children can be in different queues #### Accessing Children Results from Parent Processor ```typescript // ✅ Good -- Parent processor reads children values const videoWorker = new Worker( "videos", async (job: Job) => { // Returns object keyed by "queueName:jobId" const childrenValues = await job.getChildrenValues(); // e.g. { "media-processing:1": { audioUrl: "..." }, "media-processing:2": { thumbnailUrl: "..." } } const results = Object.values(childrenValues); await publishVideo(job.data.videoId, results); }, { connection: createWorkerConnection() }, ); ``` **Why good:** `getChildrenValues()` provides all child results, parent processor only runs after all children succeed **Gotcha:** `getChildrenValues()` returns an object keyed by `"queueName:jobId"`, not an array. If you need ordered results, include an index in child job data. --- ## Rate Limiting ### Global Rate Limiter The `limiter` option on Worker is global across all Worker instances on the same queue. ```typescript // ✅ Good -- Rate-limited API consumer const RATE_LIMIT_MAX = 10; const RATE_LIMIT_DURATION_MS = 1000; const apiWorker = new Worker( "api-calls", async (job: Job) => { const response = await callExternalApi(job.data); return response; }, { connection: createWorkerConnection(), limiter: { max: RATE_LIMIT_MAX, duration: RATE_LIMIT_DURATION_MS }, }, ); ``` **Why good:** Global limit (10 jobs/second across ALL workers on this queue), named constants, prevents API throttling ### Manual / Dynamic Rate Limiting For APIs that return rate limit headers, use dynamic rate limiting. ```typescript // ✅ Good -- Dynamic rate limiting based on API response const dynamicWorker = new Worker( "api-calls", async (job: Job) => { const response = await callExternalApi(job.data); if (response.status === 429) { const retryAfterMs = Number(response.headers["retry-after"]) * 1000; await dynamicWorker.rateLimit(retryAfterMs); throw Worker.RateLimitError(); } return response.data; }, { connection: createWorkerConnection() }, ); ``` **Why good:** `worker.rateLimit(duration)` pauses the entire queue for the specified duration, `Worker.RateLimitError()` signals BullMQ to retry after the limit expires (does not count as a failed attempt) --- ## Job Schedulers (Repeatable/Cron Jobs) `upsertJobScheduler` (v5.16.0+) replaces the deprecated `repeat` option. It is idempotent -- safe to call on every application startup. ```typescript // ✅ Good -- Fixed interval scheduler const CLEANUP_INTERVAL_MS = 3_600_000; // 1 hour await emailQueue.upsertJobScheduler( "hourly-cleanup", { every: CLEANUP_INTERVAL_MS }, { name: "cleanup-expired", data: { type: "expired-sessions" }, opts: { removeOnComplete: true }, }, ); // ✅ Good -- Cron-based scheduler await emailQueue.upsertJobScheduler( "nightly-report", { pattern: "0 0 2 * * *" }, // Every day at 2:00 AM { name: "generate-report", data: { type: "daily-summary" }, opts: { attempts: 3, backoff: { type: "exponential", delay: 5000 } }, }, ); ``` **Why good:** `upsert` is idempotent (won't duplicate on restart), template defines job name/data/opts, cron via `pattern` field #### Managing Schedulers ```typescript // Remove a scheduler await emailQueue.removeJobScheduler("hourly-cleanup"); // List all active schedulers const schedulers = await emailQueue.getJobSchedulers(); for (const scheduler of schedulers) { console.log(scheduler.id, scheduler.pattern ?? scheduler.every); } ``` **Gotcha:** `every` intervals align to the clock. A 2000ms interval fires at 0s, 2s, 4s -- not 2s after you called `upsertJobScheduler`. ```typescript // ❌ Bad -- Using deprecated repeat option await emailQueue.add( "cleanup", { type: "expired" }, { repeat: { every: 3600000 }, // Deprecated since v5.16.0 }, ); ``` **Why bad:** `repeat` option is deprecated; `upsertJobScheduler` provides better management (list, remove, upsert semantics) --- ## QueueEvents for Global Monitoring QueueEvents listens to all Workers on a queue using Redis Streams. Requires its own dedicated connection. ```typescript // ✅ Good -- Global event monitoring import { QueueEvents } from "bullmq"; const QUEUE_NAME = "emails"; const queueEvents = new QueueEvents(QUEUE_NAME, { connection: createWorkerConnection(), }); queueEvents.on("completed", ({ jobId, returnvalue }) => { console.log(`Job ${jobId} completed:`, returnvalue); }); queueEvents.on("failed", ({ jobId, failedReason }) => { console.error(`Job ${jobId} failed:`, failedReason); }); queueEvents.on("progress", ({ jobId, data }) => { console.log(`Job ${jobId} progress:`, data); }); queueEvents.on("stalled", ({ jobId }) => { console.warn(`Job ${jobId} stalled -- processor may be blocked`); }); ``` **Why good:** Global monitoring across all Workers, uses Redis Streams (reliable even across reconnections), stalled event helps detect CPU-blocked processors #### Waiting for a Specific Job ```typescript // ✅ Wait for a specific job to complete (useful in HTTP handlers) const job = await emailQueue.add("send-receipt", data); const result = await job.waitUntilFinished(queueEvents); // result = return value from the processor ``` **Why good:** `waitUntilFinished` blocks until the specific job completes or fails, useful for synchronous-feeling endpoints that need the result **Gotcha:** `waitUntilFinished` will throw if the job fails. Wrap in try-catch. --- ## Concurrency ### Local Concurrency (Per Worker) ```typescript // ✅ Good -- I/O-bound worker with concurrency const CONCURRENCY = 20; const worker = new Worker( "api-calls", async (job: Job) => { return await callExternalApi(job.data); }, { connection: createWorkerConnection(), concurrency: CONCURRENCY, }, ); // Dynamic adjustment at runtime worker.concurrency = 5; ``` **Why good:** Named constant for concurrency, I/O-bound work benefits from parallel processing, dynamic adjustment for backpressure **When NOT to use concurrency:** CPU-intensive processors block the event loop, preventing BullMQ from renewing job locks. This causes stalled jobs. Use sandboxed processors instead. --- ## Sandboxed Processors Run processors in a separate process (or worker thread) to prevent CPU-intensive work from blocking the event loop. ```typescript // ✅ Good -- Sandboxed processor (separate file) import { Worker } from "bullmq"; import path from "node:path"; const processorPath = path.join(__dirname, "processors", "image-resize.js"); const imageWorker = new Worker("images", processorPath, { connection: createWorkerConnection(), useWorkerThreads: true, // Use Node.js Worker Threads (lighter than child processes) concurrency: 4, }); ``` ```typescript // processors/image-resize.ts -- the processor file import type { SandboxedJob } from "bullmq"; export default async function (job: SandboxedJob) { // CPU-intensive work here -- runs in a separate thread const result = await resizeImage(job.data.imagePath, job.data.dimensions); return { outputPath: result.path }; } ``` **Why good:** `useWorkerThreads: true` is lighter than child processes, CPU work cannot block the main event loop, lock renewal continues uninterrupted **Gotcha:** The processor file must export a default function (this is the one exception to the named-export convention). ESM projects can use `pathToFileURL()` from `node:url` instead of `path.join`. --- ## Custom Backoff Strategies Define custom retry delay logic in the Worker settings. ```typescript // ✅ Good -- Custom backoff with jitter const BASE_DELAY_MS = 1000; const MAX_DELAY_MS = 30_000; const worker = new Worker("emails", processor, { connection: createWorkerConnection(), settings: { backoffStrategy: (attemptsMade: number) => { // Exponential with jitter const exponential = Math.min( Math.pow(2, attemptsMade) * BASE_DELAY_MS, MAX_DELAY_MS, ); const jitter = Math.random() * BASE_DELAY_MS; return exponential + jitter; }, }, }); // Reference the custom strategy when adding jobs await emailQueue.add("send", data, { attempts: 5, backoff: { type: "custom" }, }); ``` **Why good:** Jitter prevents thundering herd when many jobs retry simultaneously, capped exponential prevents excessive delays, custom logic for domain-specific retry behavior **Key behavior:** A custom backoff returning `0` moves the job to the end of the waiting list. Returning `-1` moves it directly to failed (no more retries). -
core.md 7.9 KB
# BullMQ Core Patterns > Related: [advanced.md](advanced.md) -- FlowProducer, rate limiting, job scheduling, events, concurrency --- ## Connection Factory BullMQ v5 requires explicit Redis connections. Workers need `maxRetriesPerRequest: null`; producers (Queue) can use a default or limited retry count for faster failure feedback. ```typescript // ✅ Good -- Connection factory with consumer/producer separation import Redis from "ioredis"; const RETRY_DELAY_BASE_MS = 50; const RETRY_DELAY_MAX_MS = 2000; const PRODUCER_MAX_RETRIES = 3; function createWorkerConnection(): Redis { const url = process.env.REDIS_URL; if (!url) throw new Error("REDIS_URL is required"); return new Redis(url, { maxRetriesPerRequest: null, // REQUIRED for Workers retryStrategy(times) { return Math.min(times * RETRY_DELAY_BASE_MS, RETRY_DELAY_MAX_MS); }, }); } function createProducerConnection(): Redis { const url = process.env.REDIS_URL; if (!url) throw new Error("REDIS_URL is required"); return new Redis(url, { maxRetriesPerRequest: PRODUCER_MAX_RETRIES, retryStrategy(times) { return Math.min(times * RETRY_DELAY_BASE_MS, RETRY_DELAY_MAX_MS); }, }); } export { createWorkerConnection, createProducerConnection }; ``` **Why good:** Separate factories for producers (fast failure for HTTP handlers) and consumers (infinite retries for background workers), named constants, env var validation ```typescript // ❌ Bad -- No maxRetriesPerRequest, hardcoded URL, shared connection import Redis from "ioredis"; const redis = new Redis("redis://localhost:6379"); const queue = new Queue("emails", { connection: redis }); const worker = new Worker("emails", processor, { connection: redis }); ``` **Why bad:** Missing `maxRetriesPerRequest: null` causes BullMQ Worker to throw, hardcoded URL breaks in production, sharing connection between Queue and Worker can cause blocking conflicts --- ## Queue and Worker with Typed Jobs Use TypeScript generics on Queue, Worker, and Job for type-safe job data and return values. ```typescript // ✅ Good -- Typed queue and worker import { Queue, Worker, type Job } from "bullmq"; interface EmailJobData { to: string; subject: string; body: string; } type EmailJobReturn = { messageId: string }; const QUEUE_NAME = "emails"; const MAX_ATTEMPTS = 3; const BACKOFF_DELAY_MS = 1000; const KEEP_COMPLETED = 200; // Producer const emailQueue = new Queue<EmailJobData>(QUEUE_NAME, { connection: createProducerConnection(), defaultJobOptions: { attempts: MAX_ATTEMPTS, backoff: { type: "exponential", delay: BACKOFF_DELAY_MS }, removeOnComplete: { count: KEEP_COMPLETED }, removeOnFail: false, }, }); // Consumer const emailWorker = new Worker<EmailJobData, EmailJobReturn>( QUEUE_NAME, async (job: Job<EmailJobData>) => { const result = await sendEmail( job.data.to, job.data.subject, job.data.body, ); return { messageId: result.id }; }, { connection: createWorkerConnection() }, ); // Error handler -- prevents unhandled errors from crashing the process emailWorker.on("error", (err) => { console.error("Worker error:", err.message); }); export { emailQueue, emailWorker }; ``` **Why good:** Generics enforce type safety on `job.data` and return values, `defaultJobOptions` applies to all jobs, `removeOnComplete` with count prevents unbounded Redis memory, error handler prevents process crash ```typescript // ❌ Bad -- No types, no error handler, no defaultJobOptions const queue = new Queue("emails", { connection: createWorkerConnection() }); const worker = new Worker( "emails", async (job) => { await sendEmail(job.data.to, job.data.subject, job.data.body); // No return value, no error handler }, { connection: createWorkerConnection() }, ); ``` **Why bad:** No type safety on job data, no error event handler risks crash, no retry/backoff config means failures are permanent --- ## Adding Jobs ```typescript // ✅ Single job with options const DELAY_MS = 60_000; const HIGH_PRIORITY = 1; await emailQueue.add( "send-welcome", { to: "user@example.com", subject: "Welcome", body: "Hello!", }, { delay: DELAY_MS, // Wait 60s before processing priority: HIGH_PRIORITY, // 1 = highest priority jobId: "welcome-user-123", // Idempotent -- same ID won't add duplicate }, ); // ✅ Bulk add -- single Redis round-trip await emailQueue.addBulk([ { name: "send-notification", data: { to: "a@example.com", subject: "Hi", body: "..." }, }, { name: "send-notification", data: { to: "b@example.com", subject: "Hi", body: "..." }, }, { name: "send-notification", data: { to: "c@example.com", subject: "Hi", body: "..." }, }, ]); ``` **Why good:** `addBulk` sends all jobs in a single round-trip, `jobId` prevents duplicates, named constants for delay and priority **Gotcha:** Job IDs must be strings in BullMQ v5. Integer IDs throw an exception. --- ## Job Progress Report progress from within a processor so external listeners can track it. ```typescript // ✅ Good -- Progress reporting from processor const worker = new Worker<ImportJobData>( "imports", async (job: Job<ImportJobData>) => { const BATCH_SIZE = 100; const { rows } = job.data; const total = rows.length; for (let i = 0; i < total; i += BATCH_SIZE) { const batch = rows.slice(i, i + BATCH_SIZE); await processBatch(batch); const PERCENTAGE_MULTIPLIER = 100; await job.updateProgress( Math.round(((i + batch.length) / total) * PERCENTAGE_MULTIPLIER), ); } return { processed: total }; }, { connection: createWorkerConnection() }, ); // Listen for progress from outside worker.on("progress", (job, progress) => { console.log(`Job ${job.id}: ${progress}%`); }); ``` **Why good:** `updateProgress` stores progress in Redis so any listener can read it, batch processing yields control to the event loop (prevents stalls) --- ## Graceful Shutdown ```typescript // ✅ Good -- Shutdown with timeout const SHUTDOWN_TIMEOUT_MS = 30_000; async function gracefulShutdown(workers: Worker[]): Promise<void> { const closePromise = Promise.all(workers.map((w) => w.close())); const timeoutPromise = new Promise<never>((_, reject) => { setTimeout( () => reject(new Error("Shutdown timed out")), SHUTDOWN_TIMEOUT_MS, ); }); try { await Promise.race([closePromise, timeoutPromise]); } catch (err) { console.error("Forced shutdown:", err); } finally { process.exit(0); } } process.on("SIGTERM", () => gracefulShutdown([emailWorker])); process.on("SIGINT", () => gracefulShutdown([emailWorker])); ``` **Why good:** Timeout prevents indefinite hang if a processor is stuck, handles both SIGTERM and SIGINT, `close()` stops accepting new jobs and waits for in-progress ones ```typescript // ❌ Bad -- No shutdown handler // Process exits immediately on SIGTERM, jobs become stalled ``` **Why bad:** In-progress jobs are abandoned and marked as stalled, then re-processed by another Worker (duplicate processing) --- ## Queue Management ```typescript // Pause/resume a queue (affects all Workers) await emailQueue.pause(); await emailQueue.resume(); // Drain -- remove all waiting and delayed jobs (does not affect active jobs) await emailQueue.drain(); // Obliterate -- remove ALL data for this queue from Redis await emailQueue.obliterate(); // Get job counts by status const counts = await emailQueue.getJobCounts( "waiting", "active", "completed", "failed", "delayed", ); // { waiting: 5, active: 2, completed: 100, failed: 3, delayed: 1 } // Get jobs by status const failedJobs = await emailQueue.getJobs(["failed"], 0, 10); for (const job of failedJobs) { console.log(job.id, job.failedReason); await job.retry(); // Move back to waiting } ``` **Why good:** `getJobCounts` for monitoring dashboards, `getJobs` with pagination for inspection, `retry()` for manual recovery, `obliterate` for clean slate in development
-
-
reference.md 6.8 KB
# BullMQ Quick Reference ## Job Options | Option | Type | Description | | ------------------ | ----------------------------- | ------------------------------------------------------ | | `delay` | `number` | Milliseconds to wait before processing | | `priority` | `number` | 1 = highest, MAX_INT = lowest (performance cost) | | `attempts` | `number` | Max processing attempts (1 = no retries) | | `backoff` | `{ type, delay }` | Retry strategy: `"fixed"`, `"exponential"`, `"custom"` | | `removeOnComplete` | `boolean \| { count?, age? }` | Auto-remove completed jobs | | `removeOnFail` | `boolean \| { count?, age? }` | Auto-remove failed jobs | | `lifo` | `boolean` | Last-in-first-out processing | | `jobId` | `string` | Custom ID (must be string, not integer) | | `timeout` | `number` | Milliseconds before job is considered timed out | ## Worker Options | Option | Type | Description | | ------------------ | ---------------------------- | ------------------------------------------------------------------ | | `connection` | `Redis \| ConnectionOptions` | **Required.** ioredis connection with `maxRetriesPerRequest: null` | | `concurrency` | `number` | Max parallel jobs per Worker instance (default: 1) | | `limiter` | `{ max, duration }` | Global rate limit across all Workers | | `autorun` | `boolean` | Start processing immediately (default: true) | | `lockDuration` | `number` | Lock TTL in ms (default: 30000) | | `stalledInterval` | `number` | Stall check interval in ms (default: 30000) | | `maxStalledCount` | `number` | Stalls before moving to failed (default: 1) | | `useWorkerThreads` | `boolean` | Use worker threads for sandboxed processors | ## Worker Events | Event | Payload | When | | ----------- | --------------------------- | --------------------------------- | | `completed` | `(job, returnvalue)` | Job finished successfully | | `failed` | `(job \| undefined, error)` | Job threw an error | | `progress` | `(job, progress)` | `job.updateProgress()` called | | `error` | `(error)` | Worker-level error (must handle) | | `drained` | `()` | Queue is empty, no more jobs | | `stalled` | `(jobId)` | Job lock expired, will be retried | ## QueueEvents Events | Event | Payload | When | | ----------- | ------------------------- | ------------------------------------ | | `completed` | `{ jobId, returnvalue }` | Any job completed across all Workers | | `failed` | `{ jobId, failedReason }` | Any job failed across all Workers | | `progress` | `{ jobId, data }` | Any job progress update | | `stalled` | `{ jobId }` | Any job stalled | | `waiting` | `{ jobId }` | Job added to the queue | | `delayed` | `{ jobId, delay }` | Job delayed | ## Job Lifecycle ``` added -> waiting -> active -> completed | +-> failed (retries exhausted) | +-> stalled -> waiting (re-queued) delayed -> waiting (after delay expires) ``` ## Decision Table | Scenario | Pattern | | ------------------------ | -------------------------------------------------- | | One-time background job | `queue.add(name, data)` | | Delayed job | `queue.add(name, data, { delay })` | | Priority job | `queue.add(name, data, { priority: 1 })` | | Bulk jobs | `queue.addBulk([...])` | | Recurring interval | `queue.upsertJobScheduler(id, { every })` | | Recurring cron | `queue.upsertJobScheduler(id, { pattern })` | | Parent-child workflow | `flowProducer.add({ children: [...] })` | | Rate-limited processing | Worker `limiter: { max, duration }` | | CPU-intensive processing | Sandboxed processor (file path) | | Dynamic rate limiting | `worker.rateLimit(ms)` + `Worker.RateLimitError()` | ## Anti-Patterns | Anti-Pattern | Fix | | ------------------------------------------- | ------------------------------------------------- | | No `connection` on constructor | Always pass `connection` (required in v5) | | Missing `maxRetriesPerRequest: null` | Use connection factory for Workers | | Integer job IDs | Use string IDs only | | No `worker.on("error", ...)` | Always register error handler | | No graceful shutdown | Handle SIGTERM/SIGINT with `worker.close()` | | No `removeOnComplete`/`removeOnFail` | Configure auto-removal to prevent memory growth | | Using `repeat` on `queue.add()` | Use `upsertJobScheduler` (deprecated since v5.16) | | Using QueueScheduler class | Removed in v4; Workers handle this now | | Sharing connection for Worker + QueueEvents | Each needs its own connection | | CPU work without sandboxed processor | Causes stalled jobs from blocked event loop | | Non-idempotent processors | BullMQ is at-least-once; design for duplicates | ## Production Checklist - [ ] Redis `maxmemory-policy` set to `noeviction` - [ ] All Workers use `maxRetriesPerRequest: null` - [ ] Error event handler on every Worker - [ ] Graceful shutdown on SIGTERM/SIGINT with timeout - [ ] `removeOnComplete`/`removeOnFail` configured - [ ] Processors are idempotent (at-least-once delivery) - [ ] CPU-intensive work uses sandboxed processors - [ ] QueueEvents has its own connection - [ ] Job IDs are strings (not integers) - [ ] Named constants for all delays, limits, and retry counts -
SKILL.md 17.4 KB
--- name: api-queue-bullmq description: Job queues, background processing, and task scheduling with BullMQ v5 --- # BullMQ Patterns > **Quick Guide:** Use BullMQ (v5.x) for background job processing, task scheduling, and workflow orchestration on top of Redis. Core classes: `Queue` (adds jobs), `Worker` (processes jobs), `QueueEvents` (global event listener), `FlowProducer` (parent-child job trees). Always pass a `connection` object to every constructor (required in v5). Set `maxRetriesPerRequest: null` on ioredis connections for Workers. Use `upsertJobScheduler` for repeatable/cron jobs (replaces deprecated repeatable API). QueueScheduler was removed in v4 -- its responsibilities are now handled by Workers automatically. --- <critical_requirements> ## CRITICAL: Before Using This Skill > **All code must follow project conventions in CLAUDE.md** (kebab-case, named exports, import ordering, `import type`, named constants) **(You MUST pass a `connection` object to every Queue, Worker, QueueEvents, and FlowProducer constructor -- BullMQ v5 throws if connection is missing)** **(You MUST set `maxRetriesPerRequest: null` on ioredis connections used by Workers -- BullMQ requires infinite retries and throws without this setting)** **(You MUST call `await worker.close()` on SIGTERM/SIGINT for graceful shutdown -- without it, in-progress jobs become stalled)** **(You MUST use `upsertJobScheduler` for repeatable/cron jobs -- the old `repeat` option on `queue.add` is deprecated since v5.16.0)** </critical_requirements> --- ## Examples - [Core Patterns](examples/core.md) -- Queue setup, Worker processing, job options, connection factory, graceful shutdown, typed jobs - [Advanced Patterns](examples/advanced.md) -- FlowProducer, rate limiting, job scheduling, QueueEvents, concurrency, sandboxed processors **Additional resources:** - [reference.md](reference.md) -- Decision frameworks, job option reference, anti-patterns, production checklist --- **Auto-detection:** BullMQ, bullmq, Queue, Worker, QueueEvents, FlowProducer, job queue, background job, worker process, job scheduler, upsertJobScheduler, rate limiter, job priority, job delay, sandboxed processor, repeatable job, cron job, flow producer, parent child jobs **When to use:** - Background processing (email sending, image processing, PDF generation) - Scheduled/cron jobs (nightly reports, periodic cleanup) - Workflow orchestration with parent-child job dependencies (FlowProducer) - Rate-limited API consumption (throttling outbound requests) - Priority-based job processing (urgent jobs before bulk operations) - Distributing CPU-intensive work across multiple workers or machines **Key patterns covered:** - Queue and Worker setup with typed job data and return values - Connection factory with `maxRetriesPerRequest: null` for Workers - Job options: delay, priority, attempts, backoff, removeOnComplete/Fail - Graceful shutdown with `worker.close()` on process signals - FlowProducer for parent-child job trees with dependency tracking - Job Schedulers for repeatable/cron jobs (`upsertJobScheduler`) - Rate limiting (global limiter and manual `Worker.RateLimitError`) - Concurrency control (local per-worker and global) - QueueEvents for global event monitoring across all workers - Sandboxed processors for CPU-intensive work **When NOT to use:** - Simple in-process timers or `setTimeout` (no persistence needed) - Real-time pub/sub messaging without persistence (use your pub/sub solution) - Data that must be processed synchronously within a request-response cycle - Queues that don't need persistence, retries, or scheduling --- <philosophy> ## Philosophy BullMQ is a **Redis-backed job queue** for Node.js that provides reliable background processing with at-least-once delivery guarantees. The core principle: **separate job production from job consumption** so your application stays responsive while work happens asynchronously. **Core principles:** 1. **Producers and consumers are decoupled** -- Any process can add jobs to a queue; any Worker can process them. This enables horizontal scaling by adding more Workers. 2. **Jobs are persistent** -- Jobs survive process restarts because they live in Redis. A crashed Worker's jobs are picked up by other Workers (or the same Worker after restart). 3. **At-least-once delivery** -- BullMQ guarantees every job is processed at least once. Use idempotent processors to handle the (rare) case of duplicate processing after a stall. 4. **Fail gracefully with retries** -- Configure `attempts` and `backoff` strategies so transient failures resolve automatically. Permanently failed jobs move to the failed set for inspection. 5. **Each Queue/Worker/QueueEvents needs its own connection** -- BullMQ manages connection state internally. Never share a single ioredis instance across multiple BullMQ classes (except Queues acting only as producers). </philosophy> --- <patterns> ## Core Patterns ### Pattern 1: Connection Factory BullMQ v5 requires an explicit Redis connection on every constructor. Workers need `maxRetriesPerRequest: null` so ioredis retries indefinitely instead of giving up. ```typescript import Redis from "ioredis"; function createBullMQConnection(): Redis { const url = process.env.REDIS_URL; if (!url) throw new Error("REDIS_URL is required"); return new Redis(url, { maxRetriesPerRequest: null }); } export { createBullMQConnection }; ``` **Why good:** `maxRetriesPerRequest: null` satisfies BullMQ's requirement, factory ensures consistent config, environment variable keeps credentials out of code ```typescript // Bad -- missing maxRetriesPerRequest const redis = new Redis("redis://localhost:6379"); const worker = new Worker("emails", processor, { connection: redis }); // BullMQ throws: "maxRetriesPerRequest must be null" ``` **Why bad:** BullMQ requires infinite retries on Worker connections and will throw at startup without `null` See [examples/core.md](examples/core.md) for the full connection factory with producer vs consumer separation. --- ### Pattern 2: Queue and Worker Setup Queue adds jobs; Worker processes them. Both accept TypeScript generics for type-safe job data and return values. ```typescript import { Queue, Worker, type Job } from "bullmq"; interface EmailJobData { to: string; subject: string; body: string; } const QUEUE_NAME = "emails"; const emailQueue = new Queue<EmailJobData>(QUEUE_NAME, { connection: createBullMQConnection(), }); const emailWorker = new Worker<EmailJobData>( QUEUE_NAME, async (job: Job<EmailJobData>) => { await sendEmail(job.data.to, job.data.subject, job.data.body); }, { connection: createBullMQConnection() }, ); ``` **Why good:** Generics enforce type safety on `job.data`, separate connections for Queue and Worker, named constant for queue name See [examples/core.md](examples/core.md) for full setup with events, error handling, and typed return values. --- ### Pattern 3: Job Options Control job behavior with options: delay, priority, retries, backoff, and auto-removal. ```typescript const MAX_ATTEMPTS = 5; const BACKOFF_DELAY_MS = 1000; const DELAY_MS = 60_000; const KEEP_COMPLETED_COUNT = 100; await emailQueue.add( "send-welcome", { to: "user@example.com", subject: "Welcome", body: "..." }, { delay: DELAY_MS, priority: 1, // 1 = highest, higher numbers = lower priority attempts: MAX_ATTEMPTS, backoff: { type: "exponential", delay: BACKOFF_DELAY_MS }, removeOnComplete: { count: KEEP_COMPLETED_COUNT }, removeOnFail: false, }, ); ``` **Why good:** Named constants for all numeric values, exponential backoff for transient failures, `removeOnComplete` with count prevents unbounded Redis memory, `removeOnFail: false` preserves failed jobs for debugging See [examples/core.md](examples/core.md) for `defaultJobOptions`, `addBulk`, and LIFO patterns. --- ### Pattern 4: Graceful Shutdown Call `worker.close()` on process signals to finish in-progress jobs before exiting. Without this, jobs become stalled and are re-processed by other Workers. ```typescript async function shutdown(workers: Worker[]): Promise<void> { await Promise.all(workers.map((w) => w.close())); process.exit(0); } process.on("SIGTERM", () => shutdown([emailWorker])); process.on("SIGINT", () => shutdown([emailWorker])); ``` **Why good:** `close()` stops accepting new jobs and waits for in-progress jobs to finish, `Promise.all` handles multiple workers, signal handlers cover both SIGTERM and SIGINT **Gotcha:** `worker.close()` has no built-in timeout. If a processor hangs, the shutdown hangs. Wrap with your own timeout if needed. See [examples/core.md](examples/core.md) for shutdown with timeout pattern. --- ### Pattern 5: FlowProducer (Parent-Child Jobs) FlowProducer creates job trees where parent jobs wait for all children to complete. The entire tree is added atomically. ```typescript import { FlowProducer } from "bullmq"; const flowProducer = new FlowProducer({ connection: createBullMQConnection() }); await flowProducer.add({ name: "publish-video", queueName: "videos", children: [ { name: "extract-audio", queueName: "media", data: { videoId: "abc" } }, { name: "generate-thumbnail", queueName: "media", data: { videoId: "abc" }, }, { name: "transcode", queueName: "media", data: { videoId: "abc", format: "mp4" }, }, ], }); ``` **Why good:** Parent waits until all children complete, atomic addition (all or nothing), children can be in different queues, parent processor can access children results via `job.getChildrenValues()` See [examples/advanced.md](examples/advanced.md) for accessing children values and nested flows. --- ### Pattern 6: Job Schedulers (Repeatable/Cron Jobs) Use `upsertJobScheduler` (v5.16.0+) for repeatable jobs. Replaces the deprecated `repeat` option. ```typescript const REPORT_INTERVAL_MS = 3_600_000; // 1 hour // Fixed interval await emailQueue.upsertJobScheduler( "hourly-digest", { every: REPORT_INTERVAL_MS }, { name: "send-digest", data: { type: "hourly" }, }, ); // Cron expression -- daily at 3:15 AM await emailQueue.upsertJobScheduler( "nightly-cleanup", { pattern: "0 15 3 * * *" }, { name: "cleanup", data: { type: "nightly" }, opts: { removeOnComplete: true }, }, ); ``` **Why good:** `upsertJobScheduler` is idempotent (safe to call on every startup), template defines job name/data/opts inherited by each produced job, cron syntax via `pattern` **Gotcha:** `every` intervals align to the clock, not to when you called `upsertJobScheduler`. A 2000ms interval fires at 0s, 2s, 4s, etc. See [examples/advanced.md](examples/advanced.md) for removing schedulers and listing active schedulers. --- ### Pattern 7: Rate Limiting Limit how many jobs a Worker processes per time window. The limiter is global across all Workers on the same queue. ```typescript const RATE_LIMIT_MAX = 10; const RATE_LIMIT_DURATION_MS = 1000; const worker = new Worker("api-calls", processor, { connection: createBullMQConnection(), limiter: { max: RATE_LIMIT_MAX, duration: RATE_LIMIT_DURATION_MS }, }); ``` **Why good:** Global rate limit (10 workers with this config still process max 10 jobs/second total), named constants for limits For dynamic rate limiting based on external API responses, use `worker.rateLimit(duration)` + `throw Worker.RateLimitError()`. See [examples/advanced.md](examples/advanced.md). --- ### Pattern 8: QueueEvents for Global Monitoring QueueEvents uses Redis Streams (not Pub/Sub) so events are reliable even across reconnections. Requires its own dedicated connection. ```typescript import { QueueEvents } from "bullmq"; const queueEvents = new QueueEvents(QUEUE_NAME, { connection: createBullMQConnection(), }); queueEvents.on("completed", ({ jobId, returnvalue }) => { console.log(`Job ${jobId} completed with: ${returnvalue}`); }); queueEvents.on("failed", ({ jobId, failedReason }) => { console.error(`Job ${jobId} failed: ${failedReason}`); }); ``` **Why good:** Monitors all Workers on a queue from a single listener, reliable delivery via Redis Streams, separate connection as required See [examples/advanced.md](examples/advanced.md) for progress tracking and waiting for specific job completion. </patterns> --- <performance> ## Performance Optimization - **Concurrency** -- Set `concurrency` on Worker options to process multiple jobs in parallel within a single Worker instance: `{ concurrency: 10 }`. Only effective for I/O-bound work. CPU-bound work blocks the event loop and causes stalled jobs -- use sandboxed processors instead. - **Sandboxed processors** -- Pass a file path instead of a function to the Worker constructor to run processors in a separate process (or worker thread with `useWorkerThreads: true`). Prevents CPU-intensive work from blocking lock renewal. - **addBulk** -- Use `queue.addBulk([...])` to add many jobs in a single Redis round-trip instead of calling `queue.add()` in a loop. - **removeOnComplete/Fail** -- Configure auto-removal to prevent unbounded Redis memory growth. Use `{ count: N }` to keep the last N jobs for debugging. - **Separate connections** -- Producers (Queue) can share a connection. Workers and QueueEvents each need their own connection due to blocking Redis commands. </performance> --- <decision_framework> ## Decision Framework ### When to Use BullMQ ``` Do you need background job processing? |-- NO -> Don't use BullMQ +-- YES -> Do you need persistence, retries, or scheduling? |-- NO -> Simple in-process queue or setTimeout may suffice +-- YES -> Do you need parent-child job dependencies? |-- YES -> BullMQ with FlowProducer +-- NO -> Do you need rate limiting or priority? |-- YES -> BullMQ with limiter/priority options +-- NO -> BullMQ with basic Queue + Worker ``` ### Which Job Pattern? ``` What kind of job scheduling do you need? |-- One-time delayed job -> queue.add() with delay option |-- Recurring on fixed interval -> upsertJobScheduler with every |-- Recurring on cron schedule -> upsertJobScheduler with pattern |-- Job that depends on other jobs -> FlowProducer with children |-- Bulk of independent jobs -> queue.addBulk([...]) ``` ### Concurrency Strategy ``` Is the processor CPU-intensive? |-- YES -> Use sandboxed processor (file path or useWorkerThreads) +-- NO -> Is it I/O-bound (network calls, DB queries)? |-- YES -> Set concurrency option (e.g., 10-50) +-- NO -> Default concurrency (1) is fine ``` </decision_framework> --- <red_flags> ## RED FLAGS **High Priority Issues:** - Missing `connection` on Queue/Worker/QueueEvents constructor -- BullMQ v5 throws at startup without it - Missing `maxRetriesPerRequest: null` on Worker connections -- BullMQ throws immediately - No graceful shutdown handler -- in-progress jobs become stalled on process exit and are re-processed by other Workers - Using deprecated `repeat` option on `queue.add()` instead of `upsertJobScheduler` -- deprecated since v5.16.0 - Using integer job IDs -- BullMQ v5 throws; IDs must be strings **Medium Priority Issues:** - No `removeOnComplete`/`removeOnFail` configured -- completed/failed jobs accumulate in Redis indefinitely - Sharing a single ioredis connection between Worker and QueueEvents -- both use blocking commands and will interfere - CPU-intensive processor without sandboxed mode -- blocks event loop, causes stalled jobs, prevents lock renewal - Missing error event handler on Worker -- `worker.on("error", ...)` prevents unhandled errors from crashing the process **Common Mistakes:** - Assuming `worker.close()` has a timeout -- it waits indefinitely for processors to finish; wrap with your own timeout - Using QueueScheduler class -- removed in BullMQ v4; its responsibilities are handled by Workers automatically - Expecting `limiter` to be per-Worker -- the rate limit is global across all Workers on the same queue - Not making processors idempotent -- BullMQ guarantees at-least-once delivery; a stalled job may be processed twice **Gotchas & Edge Cases:** - `upsertJobScheduler` `every` intervals align to the clock (0s, 2s, 4s), not to when you called the method - `worker.close()` does not cancel running processors -- it waits for them to finish naturally - Job priority has a performance cost -- BullMQ uses a different data structure for priority queues; skip if not needed - `QueueEvents` uses Redis Streams internally -- ensure your Redis instance has sufficient memory for stream data - FlowProducer adds the entire tree atomically -- if any child fails validation, none are added - `job.getChildrenValues()` returns an object keyed by `"queueName:jobId"` -- not an array - Redis must have `maxmemory-policy` set to `noeviction` -- BullMQ relies on keys not being evicted </red_flags> --- <critical_reminders> ## CRITICAL REMINDERS > **All code must follow project conventions in CLAUDE.md** (kebab-case, named exports, import ordering, `import type`, named constants) **(You MUST pass a `connection` object to every Queue, Worker, QueueEvents, and FlowProducer constructor -- BullMQ v5 throws if connection is missing)** **(You MUST set `maxRetriesPerRequest: null` on ioredis connections used by Workers -- BullMQ requires infinite retries and throws without this setting)** **(You MUST call `await worker.close()` on SIGTERM/SIGINT for graceful shutdown -- without it, in-progress jobs become stalled)** **(You MUST use `upsertJobScheduler` for repeatable/cron jobs -- the old `repeat` option on `queue.add` is deprecated since v5.16.0)** **Failure to follow these rules will cause startup crashes, stalled jobs, and unreliable job processing.** </critical_reminders>
Comments (0)
Sign in to join the conversation.
Reviews (0)
No reviews yet.
No comments yet.