Job Queues
BullMQ (Node.js — recommended)
import { Queue, Worker, QueueEvents } from 'bullmq';
import IORedis from 'ioredis';
const connection = new IORedis(process.env.REDIS_URL);
// Define queue
const emailQueue = new Queue('emails', { connection });
// Add job
await emailQueue.add('send-welcome', {
userId: '123',
template: 'welcome',
}, {
attempts: 3,
backoff: { type: 'exponential', delay: 2000 },
removeOnComplete: { age: 24 * 3600 }, // Cleanup after 24h
removeOnFail: { age: 7 * 24 * 3600 },
});
// Add delayed job
await emailQueue.add('send-reminder', { userId: '123' }, {
delay: 24 * 60 * 60 * 1000, // 24 hours
});
// Add prioritized job
await emailQueue.add('send-alert', { orderId: '456' }, {
priority: 1, // Lower number = higher priority
});
Worker
const worker = new Worker('emails', async (job) => {
switch (job.name) {
case 'send-welcome':
await sendWelcomeEmail(job.data.userId, job.data.template);
break;
case 'send-reminder':
await sendReminderEmail(job.data.userId);
break;
}
// Report progress
await job.updateProgress(50);
await doMoreWork();
await job.updateProgress(100);
}, {
connection,
concurrency: 5,
limiter: { max: 10, duration: 1000 }, // Rate limit: 10 jobs/sec
});
worker.on('completed', (job) => console.log(`Job ${job.id} completed`));
worker.on('failed', (job, err) => console.error(`Job ${job?.id} failed:`, err.message));
Celery (Python)
from celery import Celery
app = Celery('tasks', broker='redis://localhost:6379/0')
@app.task(bind=True, max_retries=3, default_retry_delay=60)
def send_email(self, user_id: str, template: str):
try:
user = get_user(user_id)
mailer.send(user.email, template)
except ConnectionError as exc:
self.retry(exc=exc)
# Dispatch
send_email.delay('user-123', 'welcome')
send_email.apply_async(args=['user-123', 'welcome'], countdown=3600) # Delay 1h
# Chain tasks
from celery import chain
workflow = chain(
process_order.s(order_id),
send_confirmation.s(),
update_inventory.s(),
)
workflow.apply_async()
Job Patterns
| Pattern |
Use Case |
| Fire-and-forget |
Email sending, notifications |
| Delayed jobs |
Reminders, scheduled tasks |
| Job chaining |
Multi-step workflows |
| Rate-limited |
External API calls |
| Priority queues |
Urgent vs batch processing |
| Unique jobs |
Prevent duplicate processing |
Monitoring (BullMQ)
const queueEvents = new QueueEvents('emails', { connection });
queueEvents.on('completed', ({ jobId, returnvalue }) => {
metrics.increment('jobs.completed', { queue: 'emails' });
});
queueEvents.on('failed', ({ jobId, failedReason }) => {
metrics.increment('jobs.failed', { queue: 'emails' });
});
// Bull Board (dashboard UI)
import { createBullBoard } from '@bull-board/api';
import { BullMQAdapter } from '@bull-board/api/bullMQAdapter';
import { ExpressAdapter } from '@bull-board/express';
const serverAdapter = new ExpressAdapter();
createBullBoard({ queues: [new BullMQAdapter(emailQueue)], serverAdapter });
app.use('/admin/queues', serverAdapter.getRouter());
Anti-Patterns
| Anti-Pattern |
Fix |
| No retry configuration |
Set attempts and backoff strategy |
| No dead letter handling |
Monitor failed jobs, set up alerts |
| Processing in request handler |
Offload to queue, return 202 Accepted |
| No concurrency limits |
Set worker concurrency and limiter |
| No job cleanup |
Configure removeOnComplete and removeOnFail |
Production Checklist
1---2name: job-queues3description: Background job queue systems. BullMQ (Node.js), Celery (Python), Sidekiq (Ruby), Spring Batch. Job scheduling, retries, priorities, concurrency control, and dead letter queues. USE WHEN: user mentions "job queue", "background job", "BullMQ", "Bull", "Celery", "worker", "async task", "task queue", "Sidekiq" DO NOT USE FOR: cron scheduling without queue - use `cron-scheduling`; message brokers - use messaging skills (Kafka, RabbitMQ)4---5# Job Queues67## BullMQ (Node.js — recommended)89```typescript10import { Queue, Worker, QueueEvents } from 'bullmq';11import IORedis from 'ioredis';1213const connection = new IORedis(process.env.REDIS_URL);1415// Define queue16const emailQueue = new Queue('emails', { connection });1718// Add job19await emailQueue.add('send-welcome', {20 userId: '123',21 template: 'welcome',22}, {23 attempts: 3,24 backoff: { type: 'exponential', delay: 2000 },25 removeOnComplete: { age: 24 * 3600 }, // Cleanup after 24h26 removeOnFail: { age: 7 * 24 * 3600 },27});2829// Add delayed job30await emailQueue.add('send-reminder', { userId: '123' }, {31 delay: 24 * 60 * 60 * 1000, // 24 hours32});3334// Add prioritized job35await emailQueue.add('send-alert', { orderId: '456' }, {36 priority: 1, // Lower number = higher priority37});38```3940### Worker4142```typescript43const worker = new Worker('emails', async (job) => {44 switch (job.name) {45 case 'send-welcome':46 await sendWelcomeEmail(job.data.userId, job.data.template);47 break;48 case 'send-reminder':49 await sendReminderEmail(job.data.userId);50 break;51 }5253 // Report progress54 await job.updateProgress(50);55 await doMoreWork();56 await job.updateProgress(100);57}, {58 connection,59 concurrency: 5,60 limiter: { max: 10, duration: 1000 }, // Rate limit: 10 jobs/sec61});6263worker.on('completed', (job) => console.log(`Job ${job.id} completed`));64worker.on('failed', (job, err) => console.error(`Job ${job?.id} failed:`, err.message));65```6667## Celery (Python)6869```python70from celery import Celery7172app = Celery('tasks', broker='redis://localhost:6379/0')7374@app.task(bind=True, max_retries=3, default_retry_delay=60)75def send_email(self, user_id: str, template: str):76 try:77 user = get_user(user_id)78 mailer.send(user.email, template)79 except ConnectionError as exc:80 self.retry(exc=exc)8182# Dispatch83send_email.delay('user-123', 'welcome')84send_email.apply_async(args=['user-123', 'welcome'], countdown=3600) # Delay 1h8586# Chain tasks87from celery import chain88workflow = chain(89 process_order.s(order_id),90 send_confirmation.s(),91 update_inventory.s(),92)93workflow.apply_async()94```9596## Job Patterns9798| Pattern | Use Case |99|---------|----------|100| Fire-and-forget | Email sending, notifications |101| Delayed jobs | Reminders, scheduled tasks |102| Job chaining | Multi-step workflows |103| Rate-limited | External API calls |104| Priority queues | Urgent vs batch processing |105| Unique jobs | Prevent duplicate processing |106107## Monitoring (BullMQ)108109```typescript110const queueEvents = new QueueEvents('emails', { connection });111112queueEvents.on('completed', ({ jobId, returnvalue }) => {113 metrics.increment('jobs.completed', { queue: 'emails' });114});115116queueEvents.on('failed', ({ jobId, failedReason }) => {117 metrics.increment('jobs.failed', { queue: 'emails' });118});119120// Bull Board (dashboard UI)121import { createBullBoard } from '@bull-board/api';122import { BullMQAdapter } from '@bull-board/api/bullMQAdapter';123import { ExpressAdapter } from '@bull-board/express';124125const serverAdapter = new ExpressAdapter();126createBullBoard({ queues: [new BullMQAdapter(emailQueue)], serverAdapter });127app.use('/admin/queues', serverAdapter.getRouter());128```129130## Anti-Patterns131132| Anti-Pattern | Fix |133|--------------|-----|134| No retry configuration | Set `attempts` and `backoff` strategy |135| No dead letter handling | Monitor failed jobs, set up alerts |136| Processing in request handler | Offload to queue, return 202 Accepted |137| No concurrency limits | Set worker `concurrency` and `limiter` |138| No job cleanup | Configure `removeOnComplete` and `removeOnFail` |139140## Production Checklist141142- [ ] Retry with exponential backoff configured143- [ ] Dead letter queue monitoring and alerts144- [ ] Worker concurrency tuned to resource limits145- [ ] Job progress tracking for long-running tasks146- [ ] Dashboard UI for job monitoring (Bull Board)147- [ ] Graceful shutdown: process in-flight jobs before exit