본문으로 건너뛰기

안동민 개발노트

본문 시작

Producer, Worker와 작업 상태

직렬화 가능한 작업 계약을 만들고 Producer·Worker·상태 조회 API를 연결해 보고서 작업의 진행과 결과를 관찰합니다.

지난 절에서는 reports 큐와 Redis 연결을 준비했습니다.

이제 보고서 생성 요청을 큐에 넣는 Producer, 작업을 처리하는 Worker, 상태를 조회하는 API를 연결합니다.

완성된 흐름은 POST /reports가 작업을 접수하고 202 Accepted를 반환한 뒤, GET /reports/jobs/:id가 진행 상태와 결과를 보여주는 구조입니다.

Queue.add 접수 뒤 202를 반환하고 BullMQ 작업이 waiting, active, delayed 재시도, completed 또는 failed로 이동하며 상태 API가 진행률과 결과 또는 오류를 조회하는 흐름

Nest · BullMQ Job Lifecycle

접수 응답과 작업 완료를 두 계약으로 나눈다

Queue.add()가 성공하면 202 Accepted와 작업 ID를 반환합니다. 실제 성공은 Worker의 process()가 정상 반환해 completed가 된 뒤에만 확정합니다.

BullMQ 작업 상태와 재시도 전이 Redis 접수 뒤 작업은 waiting에서 active로 이동한다. process 함수가 정상 반환하면 completed와 result가 되고, Error를 던져 시도가 남으면 exponential backoff 동안 delayed에 있다가 waiting으로 돌아간다. 시도를 모두 소진하면 failed와 failedReason이 된다. ADD FETCH RESOLVE EXHAUSTED THROW · RETRY DELAY DONE Redis 접수 jobId · 202 waiting 실행 대기 active progress 0…100 completed returnvalue delayed backoff failed failedReason
  1. 접수를 확인한다

    Queue.add()가 Redis 등록에 성공한 뒤에만 202, jobId, statusUrl을 반환합니다.

  2. 상태는 조회마다 달라질 수 있다

    enqueue 직후에도 Worker와 경쟁하고 상태와 결과는 원자적으로 함께 읽히지 않으므로, 전이 경계에서는 다음 poll로 확인합니다.

  3. Worker가 실행한다

    active에서 진행률을 기록하지만 100도 완료 신호는 아닙니다.

  4. 반환하거나 재시도한다

    정상 반환은 completed와 결과가 되고, 예외는 남은 시도에 따라 delayed 뒤 재실행되거나 최종 failed가 됩니다.

HTTP 접수

202는 완료가 아니다

접수 응답은 작업 ID와 상태 URL만 확정합니다. 실제 상태는 Worker 속도에 따라 이미 바뀌었을 수 있습니다.

성공 기준

정상 반환이 완료를 만든다

progress: 100 뒤에도 실패할 수 있습니다. process()의 반환과 completed를 함께 확인합니다.

상태 조회

소유권과 보존 기간을 확인한다

state, 진행률, 시도 횟수, 결과·오류를 반환하되 다른 tenant, 알 수 없거나 제거된 작업은 같은 404로 숨깁니다.

removeOnCompleteremoveOnFail로 기록이 제거되면 결과 조회와 custom ID의 중복 방지 기간도 끝납니다. Redis 내구성은 persistence 구성에 달려 있고, lock 상실 시 at-least-once로 중복 실행될 수 있으므로 부수 효과는 멱등해야 합니다.


직렬화 가능한 작업 계약 만들기

BullMQ는 작업 데이터를 Redis에 저장합니다.

따라서 함수, 클래스 인스턴스, 스트림, HTTP 요청 객체를 payload에 넣지 않습니다.

ID와 숫자, 문자열처럼 다시 읽을 수 있는 값만 전달하고 Worker가 필요한 데이터를 저장소에서 다시 조회하게 합니다.

src/reports/report-job.types.ts
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 사이의 내부 계약을 표현합니다.

src/reports/dto/create-report.dto.ts
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에 먼저 추가합니다.

src/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에 서버가 확인한 tenantIduserId를 넣었다고 가정합니다. Controller는 그 두 문자열만 작업 데이터로 복사하고, 토큰·비밀번호·원본 요청 객체는 넘기지 않습니다.

HTTP 요청에서 인증된 식별자만 직렬화 가능한 BullMQ 작업 데이터로 만들고 Redis에 저장한 뒤 Worker가 tenant 범위에서 원본을 다시 조회하는 보안 경계

Nest · Serializable Job Contract

payload에는 재조회할 식별자만 남긴다

HTTP 입력과 인증 주체를 검증한 뒤 작은 plain object만 Redis에 저장합니다. Worker는 그 식별자로 권한 범위 안의 최신 원본을 다시 읽습니다.

1 · HTTP 경계

DTO와 인증 주체를 분리한다

DTO는 reportIdrows를 검증합니다. tenantIduserId는 요청 본문이 아니라 인증 Guard가 확인한 주체에서 가져옵니다.

2 · job.data

직렬화 가능한 최소 계약

reportId, tenantId, requestedBy, rows처럼 JSON으로 저장할 작은 값만 전달합니다.

Redis의 작업 데이터는 평문으로 남을 수 있으므로 비밀과 원본 데이터는 넣지 않습니다.

3 · Worker

서버 저장소에서 다시 조회한다

Worker는 tenantId + reportId로 권한 범위 안의 보고서 원본을 읽고 처리합니다. 오래된 요청 객체나 연결을 복원하려 하지 않습니다.

payload 금지

실행 환경과 비밀을 직렬화하지 않는다

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개씩만 남깁니다.

src/reports/reports.service.ts
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에서 다시 읽습니다.

src/reports/reports.controller.ts
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() 메서드를 구현합니다.

src/reports/reports.processor.ts
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에 머문 뒤 다시 waitingactive를 거칩니다. attempts를 모두 소진하면 최종 failed가 되고 failedReason을 조회할 수 있습니다.

처리 중에는 updateProgress()로 숫자나 직렬화 가능한 객체를 기록할 수 있습니다.

percent: 100은 Worker가 기록한 관찰 값일 뿐 완료 신호가 아닙니다. 그 다음 코드가 실패할 수도 있으므로 클라이언트는 state === 'completed'result를 기준으로 성공을 판단합니다.

모듈에 등록

Controller, Producer Service, Worker를 같은 기능 모듈에 등록합니다.

src/reports/reports.module.ts
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 {}

ReportsProcessorproviders에서 빼면 작업은 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와 경쟁해 이미 activecompleted일 수 있기 때문입니다.

응답의 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
}

실패를 구분하는 기준

관찰 결과의미확인할 코드
요청 자체가 400DTO 검증 실패Controller 진입 전 ValidationPipe
작업이 계속 waitingWorker가 없음Processor provider와 큐 이름
작업이 delayed 뒤 다시 activeWorker 예외 뒤 재시도 대기attemptsMade, attempts, backoff
작업이 failedWorker가 재시도를 모두 소진함failedReason, Worker 로그
작업 조회가 404ID가 없거나, 기록이 제거됐거나, 현재 주체가 소유하지 않음권한 검사, removeOnComplete, removeOnFail
completed인데 파일이 없음반환값과 실제 저장이 어긋남부수 효과 완료 기준
같은 파일이 두 번 생성됨retry·stalled 복구로 중복 실행됨업무 키와 멱등 저장

보존과 중복 실행의 경계

removeOnComplete: 100removeOnFail: 100은 최근 기록 수를 제한합니다. 제거된 작업은 상태·결과·오류를 더는 조회할 수 없고, 같은 custom jobId도 다시 추가할 수 있으므로 이 옵션을 영구 업무 원장이나 영구 중복 방지 장치로 사용하지 않습니다.

BullMQ는 보통 한 번 처리를 목표로 하지만 lock 상실이나 Worker 종료 같은 최악의 경우에는 at-least-once가 되어 같은 작업이 다시 실행될 수 있습니다. Redis 장애 뒤 남는 상태 또한 AOF·RDB, 복제와 백업 구성의 내구성 범위를 넘지 않습니다.

따라서 Worker는 tenantIdreportId로 권한 범위 안의 원본을 다시 조회하고, 외부 저장 같은 부수 효과는 업무 키로 멱등하게 만들어야 합니다. 상태 API는 존재하지 않는 작업, 보존 정책으로 제거된 작업, 다른 tenant의 작업에 같은 404를 반환해 작업 ID 존재 여부를 노출하지 않습니다.


Producer는 실행 방법을 알지 않고, Worker는 HTTP 응답을 알지 않습니다.

두 계층은 직렬화 가능한 작업 계약과 큐 이름으로만 연결됩니다.

다음 절에서는 재시도 대상과 즉시 실패 대상을 나누고, jitter·중복 제거·멱등 저장·수동 복구까지 운영 기준으로 확장합니다.