2026년 08월 16일

NestJS에서 BullMQ 작업 완료 상태를 관리하는 방법

BullMQ는 Redis를 기반으로 동작하는 강력한 메시지 큐 라이브러리로, NestJS 프로젝트에서 비동기 작업을 처리할 때 널리 사용됩니다. 그런데 단순히 작업을 큐에 넣고 처리하는 것을 넘어서, 작업이 완료된 이후의 상태를 어떻게 관리하느냐는 별도의 설계가 필요한 문제입니다. 이 글에서는 NestJS 환경에서 BullMQ 작업의 완료 상태를 추적하고 제어하는 실용적인 방법을 정리합니다.

BullMQ의 작업 상태 개요

BullMQ에서 작업(Job)은 생애주기 동안 여러 상태를 거칩니다. 대표적인 상태는 다음과 같습니다.

  • waiting: 큐에 추가되어 처리 대기 중
  • active: 워커가 현재 처리 중
  • completed: 성공적으로 처리 완료
  • failed: 처리 중 오류 발생
  • delayed: 지연 실행이 예약된 상태
  • prioritized: 우선순위 대기 중

완료 상태(completed)와 실패 상태(failed)는 특히 중요합니다. BullMQ는 기본적으로 완료된 작업과 실패한 작업을 Redis에 일정 개수 보관하며, 이 설정은 큐를 생성할 때 defaultJobOptions로 제어할 수 있습니다.

NestJS에서 BullMQ 기본 설정

NestJS에서 BullMQ를 사용하려면 @nestjs/bullmq 패키지를 설치하고 모듈에 등록해야 합니다.

npm install @nestjs/bullmq bullmq

BullModule을 AppModule이나 해당 기능 모듈에 등록합니다.

import { BullModule } from '@nestjs/bullmq';

@Module({
  imports: [
    BullModule.forRoot({
      connection: {
        host: 'localhost',
        port: 6379,
      },
    }),
    BullModule.registerQueue({
      name: 'task-queue',
      defaultJobOptions: {
        removeOnComplete: {
          count: 100,  // 완료된 작업 최대 100개 보관
          age: 3600,   // 1시간 경과 후 제거
        },
        removeOnFail: {
          count: 50,
        },
      },
    }),
  ],
})
export class AppModule {}

removeOnComplete와 removeOnFail 옵션을 통해 Redis에 축적되는 데이터를 관리할 수 있습니다. count는 보관할 최대 개수, age는 보관 유지 시간(초)입니다.

작업 완료 이벤트 리스닝

BullMQ는 작업 상태 변화에 따라 이벤트를 발생시킵니다. @nestjs/bullmq에서는 @OnWorkerEvent 데코레이터를 사용해 워커 내에서 이벤트를 처리할 수 있습니다.

import { Processor, WorkerHost, OnWorkerEvent } from '@nestjs/bullmq';
import { Job } from 'bullmq';

@Processor('task-queue')
export class TaskProcessor extends WorkerHost {
  async process(job: Job): Promise<any> {
    // 실제 작업 처리 로직
    const result = await this.handleTask(job.data);
    return result;
  }

  @OnWorkerEvent('completed')
  onCompleted(job: Job, result: any) {
    console.log(`Job ${job.id} completed with result:`, result);
  }

  @OnWorkerEvent('failed')
  onFailed(job: Job, error: Error) {
    console.error(`Job ${job.id} failed with error:`, error.message);
  }

  @OnWorkerEvent('progress')
  onProgress(job: Job, progress: number | object) {
    console.log(`Job ${job.id} progress:`, progress);
  }

  private async handleTask(data: any): Promise<any> {
    // 비즈니스 로직 처리
    return { status: 'done', processedAt: new Date() };
  }
}

completed 이벤트에서 반환된 결과값(result)에 접근할 수 있습니다. 이 결과는 process 메서드가 반환한 값이며, Redis에 직렬화되어 저장됩니다.

작업 상태를 외부에서 조회하기

특정 작업의 상태를 API 엔드포인트나 서비스에서 직접 조회해야 할 때는 Queue 인스턴스를 주입받아 사용합니다.

import { InjectQueue } from '@nestjs/bullmq';
import { Queue, Job } from 'bullmq';
import { Injectable } from '@nestjs/common';

@Injectable()
export class TaskService {
  constructor(
    @InjectQueue('task-queue') private readonly taskQueue: Queue,
  ) {}

  async getJobStatus(jobId: string) {
    const job = await this.taskQueue.getJob(jobId);
    if (!job) {
      return { status: 'not_found' };
    }

    const state = await job.getState();
    const result = job.returnvalue;
    const failedReason = job.failedReason;

    return {
      id: job.id,
      state,
      result,
      failedReason,
      processedOn: job.processedOn,
      finishedOn: job.finishedOn,
    };
  }

  async addTask(data: any): Promise<string> {
    const job = await this.taskQueue.add('process-task', data);
    return job.id;
  }
}

job.getState()는 현재 작업 상태를 문자열로 반환합니다. job.returnvalue는 완료된 작업의 반환값이며, job.failedReason은 실패 시 오류 메시지를 담고 있습니다.

작업 완료 상태를 별도 저장소에 기록하기

removeOnComplete: true로 설정하면 작업이 완료되는 즉시 Redis에서 삭제됩니다. 이 경우 완료된 작업의 결과를 나중에 조회할 수 없게 됩니다. 따라서 완료 이력을 지속적으로 유지해야 하는 경우에는 데이터베이스에 별도로 기록하는 패턴이 유용합니다.

@Processor('task-queue')
export class TaskProcessor extends WorkerHost {
  constructor(private readonly jobResultRepository: JobResultRepository) {
    super();
  }

  async process(job: Job): Promise<any> {
    const result = await this.handleTask(job.data);
    return result;
  }

  @OnWorkerEvent('completed')
  async onCompleted(job: Job, result: any) {
    await this.jobResultRepository.save({
      jobId: job.id,
      queueName: job.queueName,
      status: 'completed',
      result,
      completedAt: new Date(),
    });
  }

  @OnWorkerEvent('failed')
  async onFailed(job: Job, error: Error) {
    await this.jobResultRepository.save({
      jobId: job.id,
      queueName: job.queueName,
      status: 'failed',
      errorMessage: error.message,
      failedAt: new Date(),
    });
  }
}

이 패턴은 Redis에서 작업 이력이 삭제된 이후에도 데이터베이스에서 완료 여부와 결과를 조회할 수 있다는 장점이 있습니다. 특히 감사 로그나 사용자에게 처리 결과를 안내해야 하는 경우에 적합합니다.

작업 진행률 업데이트

오래 걸리는 작업의 경우 진행률을 실시간으로 업데이트하면 클라이언트에서 진행 상황을 추적하기 수월합니다. job.updateProgress() 메서드를 사용합니다.

async process(job: Job): Promise<any> {
  const items = job.data.items;
  const total = items.length;

  for (let i = 0; i < total; i++) {
    await this.processItem(items[i]);
    await job.updateProgress(Math.round(((i + 1) / total) * 100));
  }

  return { processedCount: total };
}

updateProgress에는 숫자(0~100) 또는 객체를 전달할 수 있습니다. 객체를 사용하면 현재 처리 중인 항목 정보 등 더 세밀한 진행 상태를 표현할 수 있습니다.

QueueEvents를 활용한 전역 이벤트 리스닝

워커 외부에서, 예를 들어 별도의 서비스나 게이트웨이에서 큐 이벤트를 수신하려면 QueueEvents를 활용합니다.

import { QueueEvents } from 'bullmq';
import { Injectable, OnModuleInit, OnModuleDestroy } from '@nestjs/common';

@Injectable()
export class TaskEventListener implements OnModuleInit, OnModuleDestroy {
  private queueEvents: QueueEvents;

  onModuleInit() {
    this.queueEvents = new QueueEvents('task-queue', {
      connection: { host: 'localhost', port: 6379 },
    });

    this.queueEvents.on('completed', ({ jobId, returnvalue }) => {
      console.log(`[Event] Job ${jobId} completed:`, returnvalue);
    });

    this.queueEvents.on('failed', ({ jobId, failedReason }) => {
      console.error(`[Event] Job ${jobId} failed:`, failedReason);
    });
  }

  async onModuleDestroy() {
    await this.queueEvents.close();
  }
}

QueueEvents는 워커와 독립적으로 동작하며, 여러 곳에서 동일한 큐 이벤트를 수신할 수 있습니다. WebSocket 게이트웨이와 연결해 클라이언트에게 실시간으로 완료 상태를 푸시하는 구조를 구현할 때도 유용합니다.

완료 상태 관리 시 고려할 점

실제 프로젝트에서 BullMQ 작업 완료 상태를 관리할 때 주의할 사항을 정리합니다.

  • Redis 메모리 관리: removeOnComplete와 removeOnFail 옵션을 적절히 설정하지 않으면 완료된 작업이 무한정 쌓여 Redis 메모리를 소모합니다.
  • 작업 ID 관리: 작업을 추가할 때 jobId를 명시적으로 지정하면 외부에서 상태를 조회하기 편리합니다. 단, 동일한 jobId가 이미 존재하는 경우 동작 방식을 고려해야 합니다.
  • 직렬화 주의: process 메서드의 반환값은 JSON 직렬화 가능한 형태여야 합니다. 클래스 인스턴스나 순환 참조 구조는 저장 시 문제가 생길 수 있습니다.
  • 이벤트 핸들러의 오류 처리: @OnWorkerEvent 핸들러 내부에서 발생하는 예외는 워커의 작업 처리에 영향을 줄 수 있으므로 내부에서 별도로 오류를 처리해야 합니다.

마치며

BullMQ를 NestJS에서 사용할 때 작업 완료 상태를 잘 관리하면 비동기 처리 흐름을 안정적으로 운영할 수 있습니다. @OnWorkerEvent 데코레이터로 워커 내부에서 완료 이벤트를 처리하거나, Queue 인스턴스로 외부에서 상태를 조회하거나, QueueEvents로 전역 이벤트를 수신하는 방법을 상황에 맞게 선택해 조합하면 요구사항에 유연하게 대응할 수 있습니다. 필요에 따라 데이터베이스에 완료 이력을 기록하는 패턴도 함께 적용하면 더욱 견고한 시스템을 구축할 수 있습니다.