Producer, Worker와 작업 상태
직렬화 가능한 작업 계약을 만들고 Producer·Worker·상태 조회 API를 연결해 보고서 작업의 진행과 결과를 관찰합니다.
지난 절에서는 reports 큐와 Redis 연결을 준비했습니다.
이제 보고서 생성 요청을 큐에 넣는 Producer, 작업을 처리하는 Worker, 상태를 조회하는 API를 연결합니다.
완성된 흐름은 POST /reports가 작업을 접수하고 202 Accepted를 반환한 뒤, GET /reports/jobs/:id가 진행 상태와 결과를 보여주는 구조입니다.
Nest · BullMQ Job Lifecycle
접수 응답과 작업 완료를 두 계약으로 나눈다
Queue.add()가 성공하면 202 Accepted와 작업 ID를 반환합니다. 실제 성공은 Worker의 process()가 정상 반환해 completed가 된 뒤에만 확정합니다.
-
접수를 확인한다
Queue.add()가 Redis 등록에 성공한 뒤에만202,jobId,statusUrl을 반환합니다. -
상태는 조회마다 달라질 수 있다
enqueue 직후에도 Worker와 경쟁하고 상태와 결과는 원자적으로 함께 읽히지 않으므로, 전이 경계에서는 다음 poll로 확인합니다.
-
Worker가 실행한다
active에서 진행률을 기록하지만100도 완료 신호는 아닙니다. -
반환하거나 재시도한다
정상 반환은
completed와 결과가 되고, 예외는 남은 시도에 따라delayed뒤 재실행되거나 최종failed가 됩니다.
202는 완료가 아니다
접수 응답은 작업 ID와 상태 URL만 확정합니다. 실제 상태는 Worker 속도에 따라 이미 바뀌었을 수 있습니다.
정상 반환이 완료를 만든다
progress: 100 뒤에도 실패할 수 있습니다. process()의 반환과 completed를 함께 확인합니다.
소유권과 보존 기간을 확인한다
state, 진행률, 시도 횟수, 결과·오류를 반환하되 다른 tenant, 알 수 없거나 제거된 작업은 같은 404로 숨깁니다.
removeOnComplete와 removeOnFail로 기록이 제거되면 결과 조회와 custom ID의 중복 방지 기간도 끝납니다. Redis 내구성은 persistence 구성에 달려 있고, lock 상실 시 at-least-once로 중복 실행될 수 있으므로 부수 효과는 멱등해야 합니다.
직렬화 가능한 작업 계약 만들기
BullMQ는 작업 데이터를 Redis에 저장합니다.
따라서 함수, 클래스 인스턴스, 스트림, HTTP 요청 객체를 payload에 넣지 않습니다.
ID와 숫자, 문자열처럼 다시 읽을 수 있는 값만 전달하고 Worker가 필요한 데이터를 저장소에서 다시 조회하게 합니다.
export interface GenerateReportData {
reportId: string;
tenantId: string;
requestedBy: string;
rows: number;
}
export interface GenerateReportResult {
reportId: string;
objectKey: string;
generatedRows: number;
}
export interface ReportProgress {
step: 'load' | 'render' | 'upload';
percent: number;
}
export interface ReportPrincipal {
tenantId: string;
userId: string;
}요청 DTO는 HTTP 입력을 검증하고, 작업 데이터 타입은 Producer와 Worker 사이의 내부 계약을 표현합니다.
import { IsInt, IsString, Max, Min } from 'class-validator';
export class CreateReportDto {
@IsString()
reportId: 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();외부 입력의 requestedBy를 그대로 믿으면 다른 사용자의 작업을 만들 수 있습니다.
이 예제는 앞서 만든 인증 Guard가 request.user에 서버가 확인한 tenantId와 userId를 넣었다고 가정합니다. Controller는 그 두 문자열만 작업 데이터로 복사하고, 토큰·비밀번호·원본 요청 객체는 넘기지 않습니다.
Nest · Serializable Job Contract
payload에는 재조회할 식별자만 남긴다
HTTP 입력과 인증 주체를 검증한 뒤 작은 plain object만 Redis에 저장합니다. Worker는 그 식별자로 권한 범위 안의 최신 원본을 다시 읽습니다.
DTO와 인증 주체를 분리한다
DTO는 reportId와 rows를 검증합니다. tenantId와 userId는 요청 본문이 아니라 인증 Guard가 확인한 주체에서 가져옵니다.
직렬화 가능한 최소 계약
reportId, tenantId, requestedBy, rows처럼 JSON으로 저장할 작은 값만 전달합니다.
Redis의 작업 데이터는 평문으로 남을 수 있으므로 비밀과 원본 데이터는 넣지 않습니다.
서버 저장소에서 다시 조회한다
Worker는 tenantId + reportId로 권한 범위 안의 보고서 원본을 읽고 처리합니다. 오래된 요청 객체나 연결을 복원하려 하지 않습니다.
실행 환경과 비밀을 직렬화하지 않는다
Request·Response, 함수·class instance, stream·열린 file, database handle, access token·password, 민감한 원본 전체는 payload에서 제외합니다.
작업 ID만으로 결과를 공개하지 않는다
상태 API는 job data의 tenant·owner를 현재 인증 주체와 비교합니다. 알 수 없거나 제거됐거나 권한 밖인 작업은 같은 404로 응답합니다.
payload는 실행에 필요한 원본 자체가 아니라 원본을 안전하게 다시 찾는 계약입니다. 큰 결과는 object storage나 database에 두고 job result에는 위치와 작은 요약만 남깁니다.
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,
ReportPrincipal,
} 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, principal: ReportPrincipal) {
const job = await this.reportsQueue.add(
GENERATE_REPORT_JOB,
{
reportId: dto.reportId,
tenantId: principal.tenantId,
requestedBy: principal.userId,
rows: dto.rows,
},
{
attempts: 3,
backoff: { type: 'exponential', delay: 1_000 },
removeOnComplete: 100,
removeOnFail: 100,
},
);
return {
jobId: job.id,
accepted: true,
statusUrl: `/reports/jobs/${job.id}`,
};
}
async getStatus(jobId: string, principal: ReportPrincipal) {
const job = await this.reportsQueue.getJob(jobId);
if (
!job ||
job.data.tenantId !== principal.tenantId ||
job.data.requestedBy !== principal.userId
) {
throw new NotFoundException('report job not found');
}
const state = await job.getState();
if (state === 'unknown') {
throw new NotFoundException('report job not found');
}
return {
jobId: job.id,
name: job.name,
state,
progress: job.progress,
attemptsMade: job.attemptsMade,
result: job.returnvalue ?? null,
failedReason: job.failedReason || null,
};
}
}Queue.add()가 성공했다는 것은 보고서가 완성됐다는 뜻이 아니라 Redis가 작업을 접수했다는 뜻입니다.
그러므로 Controller는 201 Created가 아니라 202 Accepted를 반환합니다.
Worker가 매우 빠르면 Queue.add()가 반환된 뒤 HTTP 응답이 도착하기 전에도 작업을 가져갈 수 있습니다. 따라서 접수 응답은 waiting을 약속하지 않고, 상태는 별도 조회 시점의 값으로 읽습니다.
getJob()과 getState()는 서로 다른 Redis 조회입니다. 그러므로 위 응답의 state, progress, result는 하나의 원자적 snapshot이 아니며, 전이 경계에서 completed인데 result가 아직 null처럼 보이면 다음 poll에서 다시 읽습니다.
import {
Body,
Controller,
Get,
HttpCode,
HttpStatus,
Param,
Post,
Req,
} from '@nestjs/common';
import { CreateReportDto } from './dto/create-report.dto';
import { ReportPrincipal } from './report-job.types';
import { ReportsService } from './reports.service';
interface AuthenticatedRequest {
user: ReportPrincipal;
}
@Controller('reports')
export class ReportsController {
constructor(private readonly reportsService: ReportsService) {}
@Post()
@HttpCode(HttpStatus.ACCEPTED)
create(@Body() dto: CreateReportDto, @Req() request: AuthenticatedRequest) {
return this.reportsService.enqueue(dto, request.user);
}
@Get('jobs/:jobId')
getStatus(
@Param('jobId') jobId: string,
@Req() request: AuthenticatedRequest,
) {
return this.reportsService.getStatus(jobId, request.user);
}
}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,
typeof GENERATE_REPORT_JOB
>,
): 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를 던졌을 때 남은 시도가 있으면 작업은 설정한 backoff 동안 delayed에 머문 뒤 다시 waiting과 active를 거칩니다. attempts를 모두 소진하면 최종 failed가 되고 failedReason을 조회할 수 있습니다.
처리 중에는 updateProgress()로 숫자나 직렬화 가능한 객체를 기록할 수 있습니다.
percent: 100은 Worker가 기록한 관찰 값일 뿐 완료 신호가 아닙니다. 그 다음 코드가 실패할 수도 있으므로 클라이언트는 state === 'completed'와 result를 기준으로 성공을 판단합니다.
모듈에 등록
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" \
-H "Authorization: Bearer <access-token>" \
-d '{"reportId":"sales-2026-07","rows":5000}'응답 예시는 다음과 같습니다.
{
"jobId": "1",
"accepted": true,
"statusUrl": "/reports/jobs/1"
}응답에 상태를 넣지 않은 이유는 enqueue 직후에도 Worker와 경쟁해 이미 active나 completed일 수 있기 때문입니다.
응답의 jobId로 여러 번 조회합니다.
curl -H "Authorization: Bearer <access-token>" \
http://localhost:3000/reports/jobs/1조회 시점에 따라 waiting, delayed, active, completed, failed 중 하나를 볼 수 있고, progress도 함께 변합니다. 이 목록은 가능한 상태이지 한 클라이언트가 모든 중간 상태를 반드시 관찰한다는 뜻은 아닙니다.
{
"jobId": "1",
"name": "report.generate",
"state": "completed",
"progress": { "step": "upload", "percent": 100 },
"attemptsMade": 1,
"result": {
"reportId": "sales-2026-07",
"objectKey": "reports/sales-2026-07.csv",
"generatedRows": 5000
},
"failedReason": null
}실패를 구분하는 기준
| 관찰 결과 | 의미 | 확인할 코드 |
|---|---|---|
| 요청 자체가 400 | DTO 검증 실패 | Controller 진입 전 ValidationPipe |
| 작업이 계속 waiting | Worker가 없음 | Processor provider와 큐 이름 |
| 작업이 delayed 뒤 다시 active | Worker 예외 뒤 재시도 대기 | attemptsMade, attempts, backoff |
| 작업이 failed | Worker가 재시도를 모두 소진함 | failedReason, Worker 로그 |
| 작업 조회가 404 | ID가 없거나, 기록이 제거됐거나, 현재 주체가 소유하지 않음 | 권한 검사, removeOnComplete, removeOnFail |
| completed인데 파일이 없음 | 반환값과 실제 저장이 어긋남 | 부수 효과 완료 기준 |
| 같은 파일이 두 번 생성됨 | retry·stalled 복구로 중복 실행됨 | 업무 키와 멱등 저장 |
보존과 중복 실행의 경계
removeOnComplete: 100과 removeOnFail: 100은 최근 기록 수를 제한합니다. 제거된 작업은 상태·결과·오류를 더는 조회할 수 없고, 같은 custom jobId도 다시 추가할 수 있으므로 이 옵션을 영구 업무 원장이나 영구 중복 방지 장치로 사용하지 않습니다.
BullMQ는 보통 한 번 처리를 목표로 하지만 lock 상실이나 Worker 종료 같은 최악의 경우에는 at-least-once가 되어 같은 작업이 다시 실행될 수 있습니다. Redis 장애 뒤 남는 상태 또한 AOF·RDB, 복제와 백업 구성의 내구성 범위를 넘지 않습니다.
따라서 Worker는 tenantId와 reportId로 권한 범위 안의 원본을 다시 조회하고, 외부 저장 같은 부수 효과는 업무 키로 멱등하게 만들어야 합니다. 상태 API는 존재하지 않는 작업, 보존 정책으로 제거된 작업, 다른 tenant의 작업에 같은 404를 반환해 작업 ID 존재 여부를 노출하지 않습니다.
Producer는 실행 방법을 알지 않고, Worker는 HTTP 응답을 알지 않습니다.
두 계층은 직렬화 가능한 작업 계약과 큐 이름으로만 연결됩니다.
다음 절에서는 재시도 대상과 즉시 실패 대상을 나누고, jitter·중복 제거·멱등 저장·수동 복구까지 운영 기준으로 확장합니다.