Queue-Based Architecture for Reliable Processing
Build reliable message queue systems with Redis, RabbitMQ, and AWS SQS. Covers dead letter queues, idempotency, and real-world processing patterns.
Introduction
Synchronous request-response isn't always the right pattern. When operations are slow, unreliable, or need to be processed in order, message queues provide reliability and scalability that direct API calls cannot match.
When to Use Queues
Good Candidates for Async Processing
// 1. Long-running operations async function generateReport(reportId: string): Promise<void> { // Instead of making the user wait... await queue.add('reports', { reportId, requestedBy: currentUser.id, requestedAt: new Date(), }); // Return immediately with a job ID return { jobId: reportId, status: 'queued' }; } // 2. External service calls that might fail async function sendNotifications(orderId: string): Promise<void> { await queue.add('notifications', { orderId, channels: ['email', 'sms', 'push'], }, { attempts: 5, backoff: { type: 'exponential', delay: 1000 }, }); } // 3. Operations that need ordering guarantees async function processPayments(customerId: string, payments: Payment[]): Promise<void> { for (const payment of payments) { await queue.add('payments', payment, { // All payments for same customer go to same worker jobId: `${customerId}-${payment.id}`, }); } }
Implementing Reliable Queue Processing
BullMQ (Redis-based) Setup
import { Queue, Worker, QueueScheduler } from 'bullmq'; import Redis from 'ioredis'; const connection = new Redis({ host: process.env.REDIS_HOST, port: 6379, maxRetriesPerRequest: null, }); // Queue for adding jobs const emailQueue = new Queue('emails', { connection }); // Scheduler for delayed jobs and retries const scheduler = new QueueScheduler('emails', { connection }); // Worker processes jobs const worker = new Worker('emails', async (job) => { const { to, subject, template, data } = job.data; try { await emailService.send({ to, subject, template, data }); return { sent: true, sentAt: new Date() }; } catch (error) { // Throwing error triggers retry throw error; } }, { connection, concurrency: 10, limiter: { max: 100, // Max 100 jobs duration: 1000, // Per second }, }); // Job lifecycle events worker.on('completed', (job, result) => { logger.info(`Job ${job.id} completed`, { result }); }); worker.on('failed', (job, error) => { logger.error(`Job ${job.id} failed`, { error: error.message }); if (job.attemptsMade >= job.opts.attempts) { // Send to dead letter queue for manual review deadLetterQueue.add('failed-emails', { originalJob: job.data, error: error.message, attempts: job.attemptsMade, }); } });
Retry Strategies
// Exponential backoff with jitter const jobOptions = { attempts: 5, backoff: { type: 'exponential', delay: 1000, }, }; // Custom backoff strategy const customBackoff = (attemptsMade: number, error: Error): number => { // Different delays based on error type if (error instanceof RateLimitError) { return error.retryAfter * 1000; } if (error instanceof ServiceUnavailableError) { // Longer delays for service outages return Math.min(60000 * Math.pow(2, attemptsMade), 3600000); } // Default exponential backoff return Math.min(1000 * Math.pow(2, attemptsMade), 30000); };
Exactly-Once Processing
Achieving exactly-once semantics requires idempotency:
class IdempotentProcessor { constructor( private redis: Redis, private ttl: number = 86400 // 24 hours ) {} async process<T>( jobId: string, processor: () => Promise<T> ): Promise<{ result: T; cached: boolean }> { const lockKey = `processing:${jobId}`; const resultKey = `result:${jobId}`; // Check if already processed const cached = await this.redis.get(resultKey); if (cached) { return { result: JSON.parse(cached), cached: true }; } // Try to acquire lock const acquired = await this.redis.set( lockKey, 'processing', 'EX', 300, // 5 minute lock 'NX' // Only if not exists ); if (!acquired) { // Another worker is processing - wait and check result await this.waitForResult(resultKey); const result = await this.redis.get(resultKey); return { result: JSON.parse(result!), cached: true }; } try { const result = await processor(); // Store result await this.redis.setex(resultKey, this.ttl, JSON.stringify(result)); return { result, cached: false }; } finally { // Release lock await this.redis.del(lockKey); } } } // Usage in worker const processor = new IdempotentProcessor(redis); worker.process('payments', async (job) => { return processor.process(job.id, async () => { // This will only execute once even if job is retried return paymentService.process(job.data); }); });
Dead Letter Queues
// Automatic DLQ routing after max retries const deadLetterQueue = new Queue('dead-letters', { connection }); // Monitor DLQ for alerts const dlqWorker = new Worker('dead-letters', async (job) => { // Log for investigation await alertService.notify({ channel: 'slack', message: `Job failed permanently: ${job.name}`, data: { originalQueue: job.data.queue, payload: job.data.payload, error: job.data.error, attempts: job.data.attempts, }, }); // Store for manual retry await failedJobsRepository.save({ jobId: job.id, ...job.data, receivedAt: new Date(), }); }, { connection }); // Manual retry from admin interface async function retryFailedJob(jobId: string): Promise<void> { const failedJob = await failedJobsRepository.findById(jobId); if (!failedJob) { throw new Error('Job not found'); } // Re-queue to original queue const queue = new Queue(failedJob.originalQueue, { connection }); await queue.add(failedJob.name, failedJob.payload, { attempts: 3, }); // Mark as retried await failedJobsRepository.markRetried(jobId); }
Conclusion
Queue-based architecture provides:
- Reliability: Jobs persist even if workers crash
- Scalability: Add workers to increase throughput
- Resilience: Automatic retries handle transient failures
- Ordering: Process jobs in sequence when needed
- Visibility: Monitor queue depths and processing times
The key is choosing the right queue technology and implementing proper idempotency for exactly-once semantics.
Related Articles
Software Architecture18 min read
Event-Driven Architecture in Enterprise Systems: Patterns and Trade-offs
A practitioner's guide to implementing event-driven architecture at scale. Covers message broker selection, event schema design, eventual consistency patterns, and lessons from production systems.
Software Architecture19 min read
Building Resilient Distributed Systems: Patterns for Fault Tolerance
Build resilient distributed systems with circuit breakers, retries, and timeouts. Production patterns for handling failures, cascading errors, and maintaining availability.
Backend Design20 min read
Database Design Patterns for Scale
Scale databases with sharding, replication, and partitioning. Covers PostgreSQL, MySQL, and MongoDB scaling patterns with real performance numbers from production systems.
Payment Integrations28 min read
Stripe Payment Integration: Production Patterns for React and Node.js
Production Stripe integration with Payment Intents, webhooks, and 3D Secure. Covers subscription billing, error handling, and PCI compliance patterns.
Security Engineering18 min read
API Security Hardening: A Practitioner's Guide
Secure your APIs with rate limiting, input validation, and CORS configuration. Production-tested checklist covering authentication, encryption, and error handling.