Backend Design

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.

Khalid Aboubakr
17 min read
Message QueueRabbitmqRedisAsync ProcessingReliabilitySqsDead Letter QueueIdempotency

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:

  1. Reliability: Jobs persist even if workers crash
  2. Scalability: Add workers to increase throughput
  3. Resilience: Automatic retries handle transient failures
  4. Ordering: Process jobs in sequence when needed
  5. 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

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.

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.