Producer, Worker와 작업 상태
직렬화 가능한 작업 계약을 만들고 Producer·Worker·상태 조회 API를 연결해 보고서 작업의 진행과 결과를 관찰합니다.
지난 절에서는 reports 큐와 Redis 연결을 준비했습니다.
이제 보고서 생성 요청을 큐에 넣는 Producer, 작업을 처리하는 Worker, 상태를 조회하는 API를 연결합니다.
완성된 흐름은 POST /reports가 작업을 접수하고 202 Accepted를 반환한 뒤, GET /reports/jobs/:id가 진행 상태와 결과를 보여주는 구조입니다.
직렬화 가능한 작업 계약 만들기
BullMQ는 작업 데이터를 Redis에 저장합니다.
따라서 함수, 클래스 인스턴스, 스트림, HTTP 요청 객체를 payload에 넣지 않습니다.
ID와 숫자, 문자열처럼 다시 읽을 수 있는 값만 전달하고 Worker가 필요한 데이터를 저장소에서 다시 조회하게 합니다.
export interface GenerateReportData {
reportId: string;
requestedBy: string;
rows: number;
}
export interface GenerateReportResult {
reportId: string;
objectKey: string;
generatedRows: number;
}
export interface ReportProgress {
step: 'load' | 'render' | 'upload';
percent: number;
}요청 DTO는 HTTP 입력을 검증하고, 작업 데이터 타입은 Producer와 Worker 사이의 내부 계약을 표현합니다.
import { IsInt, IsString, Max, Min } from 'class-validator';
export class CreateReportDto {
@IsString()
reportId: string;
@IsString()
requestedBy: string;
@IsInt()
@Min(1)
@Max(100_000)
rows: number;
}이 예제는 앞 장에서 전역 ValidationPipe를 설정했다고 가정합니다.
설정하지 않았다면 main.ts에 먼저 추가합니다.
import { ValidationPipe } from '@nestjs/common';
import { NestFactory } from '@nestjs/core';
import { AppModule } from './app.module';
async function bootstrap() {
const app = await NestFactory.create(AppModule);
app.useGlobalPipes(new ValidationPipe({ whitelist: true, transform: true }));
await app.listen(process.env.PORT ?? 3000);
}
void bootstrap();Producer에서 작업 추가하기
Producer는 큐를 주입받아 작업 이름, 데이터, 보관 정책을 전달합니다.
작업이 완료되거나 실패한 기록을 무한히 쌓지 않도록 최근 100개씩만 남깁니다.
import { InjectQueue } from '@nestjs/bullmq';
import { Injectable, NotFoundException } from '@nestjs/common';
import { Queue } from 'bullmq';
import { CreateReportDto } from './dto/create-report.dto';
import {
GenerateReportData,
GenerateReportResult,
} from './report-job.types';
import { GENERATE_REPORT_JOB, REPORT_QUEUE } from './reports.constants';
@Injectable()
export class ReportsService {
constructor(
@InjectQueue(REPORT_QUEUE)
private readonly reportsQueue: Queue<
GenerateReportData,
GenerateReportResult,
typeof GENERATE_REPORT_JOB
>,
) {}
async enqueue(dto: CreateReportDto) {
const job = await this.reportsQueue.add(
GENERATE_REPORT_JOB,
{
reportId: dto.reportId,
requestedBy: dto.requestedBy,
rows: dto.rows,
},
{
removeOnComplete: 100,
removeOnFail: 100,
},
);
return {
jobId: job.id,
state: await job.getState(),
statusUrl: `/reports/jobs/${job.id}`,
};
}
async getStatus(jobId: string) {
const job = await this.reportsQueue.getJob(jobId);
if (!job) {
throw new NotFoundException(`report job ${jobId} not found`);
}
return {
jobId: job.id,
name: job.name,
state: await job.getState(),
progress: job.progress,
attemptsMade: job.attemptsMade,
result: job.returnvalue ?? null,
failedReason: job.failedReason || null,
};
}
}Queue.add()가 성공했다는 것은 보고서가 완성됐다는 뜻이 아니라 Redis가 작업을 접수했다는 뜻입니다.
그러므로 Controller는 201 Created가 아니라 202 Accepted를 반환합니다.
import {
Body,
Controller,
Get,
HttpCode,
HttpStatus,
Param,
Post,
} from '@nestjs/common';
import { CreateReportDto } from './dto/create-report.dto';
import { ReportsService } from './reports.service';
@Controller('reports')
export class ReportsController {
constructor(private readonly reportsService: ReportsService) {}
@Post()
@HttpCode(HttpStatus.ACCEPTED)
create(@Body() dto: CreateReportDto) {
return this.reportsService.enqueue(dto);
}
@Get('jobs/:jobId')
getStatus(@Param('jobId') jobId: string) {
return this.reportsService.getStatus(jobId);
}
}Worker에서 작업 처리하기
BullMQ용 Nest Worker는 WorkerHost를 상속하고 큐마다 하나의 process() 메서드를 구현합니다.
import { Processor, WorkerHost } from '@nestjs/bullmq';
import { Job } from 'bullmq';
import {
GenerateReportData,
GenerateReportResult,
} from './report-job.types';
import { GENERATE_REPORT_JOB, REPORT_QUEUE } from './reports.constants';
const wait = (milliseconds: number) =>
new Promise((resolve) => setTimeout(resolve, milliseconds));
@Processor(REPORT_QUEUE)
export class ReportsProcessor extends WorkerHost {
async process(
job: Job<GenerateReportData, GenerateReportResult, string>,
): Promise<GenerateReportResult> {
if (job.name !== GENERATE_REPORT_JOB) {
throw new Error(`unsupported report job: ${job.name}`);
}
await job.updateProgress({ step: 'load', percent: 20 });
await wait(700);
await job.updateProgress({ step: 'render', percent: 70 });
await wait(700);
await job.updateProgress({ step: 'upload', percent: 95 });
await wait(700);
await job.updateProgress({ step: 'upload', percent: 100 });
return {
reportId: job.data.reportId,
objectKey: `reports/${job.data.reportId}.csv`,
generatedRows: job.data.rows,
};
}
}process()가 값을 반환하면 작업은 completed 상태가 되고 반환값은 job.returnvalue에 저장됩니다.
Error를 던지면 failed 상태로 이동합니다.
처리 중에는 updateProgress()로 숫자나 직렬화 가능한 객체를 기록할 수 있습니다.
모듈에 등록
Controller, Producer Service, Worker를 같은 기능 모듈에 등록합니다.
import { BullModule } from '@nestjs/bullmq';
import { Module } from '@nestjs/common';
import { ReportsController } from './reports.controller';
import { ReportsProcessor } from './reports.processor';
import { ReportsService } from './reports.service';
import { REPORT_QUEUE } from './reports.constants';
@Module({
imports: [BullModule.registerQueue({ name: REPORT_QUEUE })],
controllers: [ReportsController],
providers: [ReportsService, ReportsProcessor],
})
export class ReportsModule {}ReportsProcessor를 providers에서 빼면 작업은 waiting에 계속 남습니다.
큐 등록만으로 Worker가 자동 생성되는 것은 아닙니다.
실행하며 상태 관찰하기
Redis와 Nest 애플리케이션을 실행한 상태에서 작업을 제출합니다.
curl -i -X POST http://localhost:3000/reports \
-H "Content-Type: application/json" \
-d '{"reportId":"sales-2026-07","requestedBy":"user-42","rows":5000}'응답 예시는 다음과 같습니다.
{
"jobId": "1",
"state": "waiting",
"statusUrl": "/reports/jobs/1"
}응답의 jobId로 여러 번 조회합니다.
curl http://localhost:3000/reports/jobs/1처리 시점에 따라 waiting → active → completed로 바뀌고, progress도 함께 변합니다.
{
"jobId": "1",
"name": "report.generate",
"state": "completed",
"progress": { "step": "upload", "percent": 100 },
"attemptsMade": 0,
"result": {
"reportId": "sales-2026-07",
"objectKey": "reports/sales-2026-07.csv",
"generatedRows": 5000
},
"failedReason": null
}실패를 구분하는 기준
| 관찰 결과 | 의미 | 확인할 코드 |
|---|---|---|
| 요청 자체가 400 | DTO 검증 실패 | Controller 진입 전 ValidationPipe |
| 작업이 계속 waiting | Worker가 없음 | Processor provider와 큐 이름 |
| 작업이 failed | Worker가 예외를 던짐 | failedReason, Worker 로그 |
| 작업 조회가 404 | 보관 한도를 넘었거나 ID가 틀림 | removeOnComplete, removeOnFail |
| completed인데 파일이 없음 | 반환값과 실제 저장이 어긋남 | 부수 효과 완료 기준 |
Producer는 실행 방법을 알지 않고, Worker는 HTTP 응답을 알지 않습니다.
두 계층은 직렬화 가능한 작업 계약과 큐 이름으로만 연결됩니다.
다음 절에서는 일시적 장애를 자동 재시도하고, 중복 요청과 중복 실행이 업무 결과를 망치지 않도록 복구 기준을 추가합니다.