Consumers

⭐ Interview Importance: HIGH
⏱️ Revision Time: 9 min

In the context of Job Queues, a Consumer (often called a Worker or Processor) is the background process that pulls jobs off the queue, executes the heavy logic, and marks the job as complete or failed.

Overview

If the Producer is the “Waitress” taking the order, the Consumer is the “Chef” in the kitchen cooking the food.

In NestJS, Consumers are classes decorated with @Processor(). They run asynchronously, totally detached from the HTTP request cycle. They maintain a persistent connection to Redis, waiting for new jobs to arrive in their designated queue.

Key Concepts

  • @Processor('queue_name'): The class-level decorator that designates a class as a Consumer for a specific queue.
  • @Process('job_name'): The method-level decorator that defines the specific function to run when a certain type of job is pulled from the queue.
  • Concurrency: A Consumer can process multiple jobs simultaneously. You configure this by setting the concurrency option on the @Process() decorator.

Code Examples

1. A Standard Consumer (Bull)

This consumer listens to the emails queue and handles two different types of jobs.

import { Processor, Process } from '@nestjs/bull';
import { Job } from 'bull';
import { EmailService } from './email.service';

@Processor('emails')
export class EmailConsumer {
  constructor(private emailService: EmailService) {}

  // Handles jobs named 'welcome_email'
  // concurrency: 5 means this method can handle 5 jobs at the exact same time
  @Process({ name: 'welcome_email', concurrency: 5 })
  async handleWelcomeEmail(job: Job<{ email: string; name: string }>) {
    console.log(`Sending welcome email to ${job.data.email}`);
    
    // If this throws an error, Bull will automatically mark the job as failed
    // and attempt to retry it (if retries were configured by the Producer).
    await this.emailService.sendWelcome(job.data.email, job.data.name);
    
    // Return value is saved in Redis
    return { status: 'Sent' };
  }

  // Handles jobs named 'password_reset'
  @Process('password_reset')
  async handlePasswordReset(job: Job<{ email: string; token: string }>) {
    await this.emailService.sendResetToken(job.data.email, job.data.token);
  }
}

2. Handling Queue Events

Consumers can listen to lifecycle events emitted by the queue to track progress or handle failures.

import { Processor, Process, OnQueueActive, OnQueueFailed } from '@nestjs/bull';
import { Job } from 'bull';

@Processor('emails')
export class EmailConsumer {
  
  @Process('welcome_email')
  async handle(job: Job) { /* ... */ }

  // Fires right before the job starts processing
  @OnQueueActive()
  onActive(job: Job) {
    console.log(`Processing job ${job.id} of type ${job.name}...`);
  }

  // Fires if the job throws an error and fails
  @OnQueueFailed()
  onError(job: Job, error: any) {
    console.error(`Job ${job.id} failed: ${error.message}`);
    // You could send a Slack alert here for critical failures!
  }
}

Best Practices

  • Idempotency is Critical: If a consumer crashes after charging a credit card, but before returning success to Bull, Bull will think the job failed and retry it. The consumer will charge the card again! A consumer must be Idempotent (safe to run multiple times). Always check the database first: if (order.isPaid) return;.
  • Database Transactions: If a job involves multiple database writes, wrap them all in a single SQL Transaction. If the worker crashes halfway through, the transaction rolls back, preventing corrupt, half-finished data from polluting your system.
  • CPU Bound vs. I/O Bound: If your consumer makes HTTP requests or database queries (I/O Bound), you can safely set a high concurrency (e.g., concurrency: 50) because Node.js is great at async I/O. However, if your consumer does image resizing or math calculations (CPU Bound), setting concurrency higher than 1 will freeze the Node Event Loop. For CPU bound tasks, use Bull’s Separate Process Sandboxing.