본문으로 건너뛰기

안동민 개발노트

본문 시작

재시도, 중복 방지와 실패 복구

일시적 실패와 영구 실패를 구분하고 지수 백오프·범위가 명확한 중복 제거·멱등 저장·감사 가능한 수동 복구를 연결합니다.

작업 큐를 운영하면 외부 저장소의 순간 장애, 네트워크 단절, 잘못된 입력처럼 성격이 다른 실패를 만나게 됩니다.

모든 실패를 같은 방식으로 재시도하면 복구 가능성이 없는 작업이 자원을 계속 차지하고, 중복 실행이 보고서를 여러 번 저장할 수 있습니다.

이번 절에서는 앞 절의 보고서 작업에 자동 재시도, 지수 백오프, 범위가 명확한 중복 제거, 멱등한 결과 저장, 수동 복구를 차례로 추가합니다.

아래 BullMQ 옵션과 Job.retry() 서명은 공식 소스 v5.77.6을 기준으로 검증했습니다. 이 문서 사이트 자체에는 bullmq@nestjs/bullmq가 설치되어 있지 않으므로, 실습 애플리케이션에서는 검증한 버전을 고정하고 Nest 공식 queue 문서의 @nestjs/bullmq 통합을 사용합니다.


재시도할 실패와 즉시 끝낼 실패

먼저 실패를 두 종류로 나눕니다.

  • 일시적 실패: 외부 저장소의 503 응답, 잠깐의 연결 끊김처럼 시간이 지나면 성공할 가능성이 있습니다.
  • 영구적 실패: 존재하지 않는 보고서 종류, 허용 범위를 벗어난 입력처럼 같은 데이터로 다시 실행해도 실패합니다.

BullMQ Worker가 일반 Error를 던지면 남은 attempts가 있을 때 작업을 다시 실행합니다.

UnrecoverableError를 던지면 남은 실행 기회와 관계없이 작업을 최종 failed 상태로 보냅니다.

실패 실험을 위해 앞 절의 작업 데이터에 선택 속성을 추가합니다. tenantIdrequestedBy는 계속 서버가 확인한 인증 주체에서 가져옵니다.

src/reports/report-job.types.ts
export interface GenerateReportData {
  reportId: string;
  tenantId: string;
  requestedBy: string;
  rows: number;
  failUntilAttempt?: number;
}

export interface GenerateReportResult {
  operationKey: string;
  reportId: string;
  objectKey: string;
  generatedRows: number;
  reused: boolean;
}

export interface ReportProgress {
  step: 'load' | 'render' | 'upload';
  percent: number;
}

DTO에는 장애 주입용 실험 필드만 추가합니다. 운영 API에서는 이 필드를 외부에 노출하지 않고 테스트 환경에서만 활성화합니다.

src/reports/dto/create-report.dto.ts
import { IsInt, IsOptional, IsString, Max, Min } from 'class-validator';

export class CreateReportDto {
  @IsString()
  reportId: string;

  @IsInt()
  @Min(1)
  @Max(100_000)
  rows: number;

  @IsOptional()
  @IsInt()
  @Min(0)
  @Max(5)
  failUntilAttempt?: number;
}

같은 업무를 나타내는 키 만들기

단순히 reportId만 중복 제거 키로 사용하면 tenant가 다르거나 결과를 바꾸는 입력이 다른 요청까지 하나로 합칠 수 있습니다.

이 예제는 작업 종류와 버전, 인증된 tenant와 사용자, 보고서 ID, 행 수를 모두 직렬화한 뒤 해시합니다. 실제 업무에서는 결과를 바꾸는 모든 입력을 포함하고, 비밀 값 자체는 payload나 키에 넣지 않습니다.

src/reports/report-operation-key.ts
import { createHash } from 'node:crypto';
import type { GenerateReportData } from './report-job.types';

export function reportOperationKey(
  data: Pick<
    GenerateReportData,
    'tenantId' | 'requestedBy' | 'reportId' | 'rows'
  >,
) {
  return createHash('sha256')
    .update(
      JSON.stringify([
        'report.generate.v1',
        data.tenantId,
        data.requestedBy,
        data.reportId,
        data.rows,
      ]),
    )
    .digest('hex');
}

해시는 범위를 정의하는 수단이지 권한 검사를 대신하지 않습니다.

Producer의 재시도 옵션

앞 절의 ReportsService.enqueue()에서 인증 주체로 작업 데이터를 만들고 옵션을 확장합니다.

src/reports/reports.service.ts (enqueue 일부)
const data: GenerateReportData = {
  reportId: dto.reportId,
  tenantId: principal.tenantId,
  requestedBy: principal.userId,
  rows: dto.rows,
  failUntilAttempt: dto.failUntilAttempt,
};

const operationKey = reportOperationKey(data);
const job = await this.reportsQueue.add(
  GENERATE_REPORT_JOB,
  data,
  {
    attempts: 4,
    backoff: {
      type: 'exponential',
      delay: 1_000,
      jitter: 0.5,
    },
    deduplication: {
      id: operationKey,
    },
    removeOnComplete: 100,
    removeOnFail: 100,
  },
);

return {
  jobId: job.id,
  accepted: true,
  statusUrl: '/reports/jobs/' + job.id,
};

attempts: 4는 최초 실행을 포함해 최대 네 번 실행할 수 있다는 뜻입니다.

delay: 1_000인 지수 백오프의 최대 대기는 첫 실패 뒤 1초, 다음 실패 뒤 2초, 그다음 실패 뒤 4초입니다. jitter: 0.5는 각 최대치의 50%에서 100% 사이로 실제 대기를 무작위화해 여러 작업이 같은 순간에 다시 몰리는 현상을 줄입니다.

재시도 간격은 외부 시스템의 복구 시간과 업무 마감 시간을 기준으로 정합니다. 짧으면 장애 중인 시스템을 재압박하고, 길면 사용자의 복구 대기가 늘어납니다.

BullMQ에서 일시적 실패는 attempts와 지수 백오프, jitter에 따라 다시 실행하고 영구 실패는 UnrecoverableError로 즉시 종료하며, 정상 반환만 completed가 되는 재시도 정책

Nest · BullMQ Retry Policy

재시도는 실패 원인과 남은 실행 기회로 결정한다

attempts: 4는 최초 실행을 포함한 네 번의 기회입니다. 일반 Error는 남은 기회와 backoff를 따르고, UnrecoverableError는 즉시 최종 failed로 갑니다.

일시적 실패

같은 입력이 나중에는 성공할 수 있음

네트워크 단절이나 외부 저장소의 503처럼 원인이 시간에 따라 바뀝니다. 일반 Error를 던져 자동 재시도 정책에 맡깁니다.

영구 실패

같은 입력으로는 결과가 바뀌지 않음

지원하지 않는 job name이나 유효하지 않은 업무 입력은 UnrecoverableError로 추가 자동 재시도 없이 종료합니다.

  1. 실행 기회 1 · 최초 실행

    정상 반환하면 completed입니다. 일반 Errorjitter: 0.5가 적용된 약 0.5–1초 대기 뒤 다음 기회로 갑니다.

  2. 실행 기회 2 · 첫 자동 재시도

    다시 실패하면 지수 백오프 최대치가 2초로 늘어 실제 대기는 약 1–2초 범위가 됩니다.

  3. 실행 기회 3 · 두 번째 자동 재시도

    다시 실패하면 최대치 4초의 약 2–4초 범위에서 대기합니다. 실제 상태는 delayed, waiting, active 사이에서 빠르게 바뀔 수 있습니다.

  4. 실행 기회 4 · 마지막 실행

    정상 반환하면 completed, 일반 Error가 계속되면 실행 예산을 소진해 최종 failed가 됩니다.

성공 경계

process() 정상 반환

progress: 100만으로는 완료가 아닙니다. 멱등 저장까지 성공한 뒤 Worker가 값을 반환해야 completed와 결과가 기록됩니다.

실패 경계

영구 실패 또는 실행 예산 소진

실패 evidence는 보존 정책 안에서만 Redis에 남습니다. 수동 retry 전에 오류와 counter를 외부 감사 저장소에 먼저 기록합니다.

BullMQ v5.77.6 기준입니다. jitter: 0.5는 각 지수 백오프 최대치의 50%에서 100% 사이를 무작위로 선택하며, 재시도는 장애를 고치는 기능이 아니라 복구될 시간을 주는 정책입니다.


중복 제거와 멱등성은 다른 문제다

위 설정은 TTL이 없는 BullMQ Simple deduplication입니다.

같은 deduplication.id의 작업이 아직 completedfailed가 아니라면 뒤의 추가 요청은 새 작업으로 저장되지 않고 deduplicated 이벤트를 냅니다.

이 범위는 큐 안의 미완료 작업 하나를 억제하는 데 그칩니다.

  • 원 작업이 completed 또는 failed가 되면 Simple deduplication 키는 끝나며 같은 업무 요청을 다시 추가할 수 있습니다.
  • Worker가 lock을 갱신하지 못해 stalled가 되면 작업은 waiting으로 돌아가 다시 실행될 수 있습니다.
  • 파일 저장 직후 Worker가 종료되면 BullMQ가 완료를 기록하지 못한 채 같은 부수 효과를 다시 수행할 수 있습니다.
  • removeOnCompleteremoveOnFail로 작업이 제거되면 큐 기록은 상태 조회나 영구 중복 방지 원장이 될 수 없습니다.

따라서 큐 중복 제거는 exactly-once 실행이나 결과 한 번 저장을 보장하지 않습니다. 최종 부수 효과를 업무 키로 멱등하게 만들어야 합니다.

실습에서는 동작을 눈으로 확인하기 위해 메모리 저장소를 사용합니다.

src/reports/report-artifacts.service.ts
import { Injectable } from '@nestjs/common';
import { GenerateReportResult } from './report-job.types';

@Injectable()
export class ReportArtifactsService {
  private readonly artifacts = new Map<string, GenerateReportResult>();

  find(operationKey: string) {
    return this.artifacts.get(operationKey);
  }

  saveOnce(operationKey: string, result: GenerateReportResult) {
    const existing = this.artifacts.get(operationKey);
    if (existing) return existing;

    this.artifacts.set(operationKey, result);
    return result;
  }
}

메모리 Map은 프로세스를 재시작하거나 Worker를 여러 대 띄우면 공유되지 않습니다.

한 프로세스 안의 예제를 보여줄 뿐이며, 운영 환경에서는 데이터베이스의 고유 제약과 원자적 upsert, 객체 저장소의 조건부 생성, 외부 API의 idempotency key로 바꿉니다.

Worker는 같은 operation key로 기존 결과를 먼저 조회하고, 정상 반환할 결과도 같은 키로 저장합니다.

src/reports/reports.processor.ts
import { Processor, WorkerHost } from '@nestjs/bullmq';
import { Job, UnrecoverableError } from 'bullmq';
import { ReportArtifactsService } from './report-artifacts.service';
import {
  GenerateReportData,
  GenerateReportResult,
} from './report-job.types';
import { reportOperationKey } from './report-operation-key';
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 {
  constructor(private readonly artifacts: ReportArtifactsService) {
    super();
  }

  async process(
    job: Job<
      GenerateReportData,
      GenerateReportResult,
      typeof GENERATE_REPORT_JOB
    >,
  ): Promise<GenerateReportResult> {
    if (job.name !== GENERATE_REPORT_JOB) {
      throw new UnrecoverableError(
        'unsupported report job: ' + job.name,
      );
    }

    if (job.data.rows < 1) {
      throw new UnrecoverableError('rows must be greater than zero');
    }

    const operationKey = reportOperationKey(job.data);
    const existing = this.artifacts.find(operationKey);
    if (existing) {
      return { ...existing, reused: true };
    }

    if (job.attemptsMade < (job.data.failUntilAttempt ?? 0)) {
      throw new Error(
        'temporary storage failure on attempt ' +
          (job.attemptsMade + 1),
      );
    }

    await job.updateProgress({ step: 'load', percent: 20 });
    await wait(400);
    await job.updateProgress({ step: 'render', percent: 70 });
    await wait(400);
    await job.updateProgress({ step: 'upload', percent: 95 });
    await wait(400);
    await job.updateProgress({ step: 'upload', percent: 100 });

    return this.artifacts.saveOnce(operationKey, {
      operationKey,
      reportId: job.data.reportId,
      objectKey:
        'reports/' + job.data.tenantId + '/' + operationKey + '.csv',
      generatedRows: job.data.rows,
      reused: false,
    });
  }
}

percent: 100은 Worker가 기록한 진행 값일 뿐 완료가 아닙니다. 그 뒤의 조건부 저장이 성공하고 process()가 정상 반환해야 BullMQ가 completedreturnvalue를 기록합니다.

기능 모듈의 provider에도 저장소를 추가합니다.

src/reports/reports.module.ts (providers 일부)
providers: [
  ReportsService,
  ReportsProcessor,
  ReportArtifactsService,
],

운영 저장소의 원자성

실제 저장소에서 find()save()를 서로 떨어진 두 쿼리로 구현하면 두 Worker가 동시에 통과할 수 있습니다.

다음 중 하나로 원자성을 확보합니다.

  • operation key에 고유 인덱스를 두고 insert 충돌 시 기존 결과를 조회합니다.
  • 데이터베이스의 INSERT ... ON CONFLICT 또는 트랜잭션을 사용합니다.
  • 객체 저장소에서 같은 key에 대한 조건부 생성 기능을 사용합니다.
  • 외부 API가 idempotency key를 지원하면 operation key를 함께 전달합니다.
인증된 업무 키, BullMQ Simple deduplication, 원자적 멱등 저장소가 서로 다른 중복 경계를 맡고, 실패 뒤에는 evidence와 권한을 확인한 후 수동 재시도하는 복구 절차

Nest · Duplicate And Recovery Boundaries

중복 억제와 결과 멱등성을 서로 다른 층에서 보장한다

deduplication.id는 미완료 작업의 뒤 추가를 억제할 뿐입니다. stalled replay와 완료·실패 뒤의 새 enqueue까지 견디려면 인증된 operation key와 원자적 sink가 필요합니다.

  1. HTTP 경계 · 인증된 operation key

    서버가 확인한 tenantId와 사용자, 작업 종류·버전, 결과를 바꾸는 입력을 묶습니다. 접수 응답은 Queue.add()가 반환한 jobIdstatusUrl을 사용합니다.

  2. Queue 경계 · Simple deduplication

    같은 키의 원 작업이 completedfailed가 되기 전까지만 뒤 추가를 무시합니다. 이 키는 exactly-once 실행이나 영구 업무 원장이 아닙니다.

  3. Sink 경계 · 원자적 멱등 저장

    unique constraint, conditional create, atomic upsert 또는 외부 idempotency key로 같은 operation 결과를 한 번만 만듭니다. lock 상실로 작업이 다시 실행돼도 기존 결과를 반환합니다.

retry 전 evidence

지워지기 전에 외부 감사 저장소에 append

jobId, 안전한 payload 식별자, failedReason, stacktrace, 두 attempt counter, 실행자·이유·시각을 먼저 남깁니다.

운영 복구 gate

권한·tenant·경쟁·원인 수정 확인

같은 operation의 새 작업이나 결과가 없는지 원자적으로 확인한 뒤 retry합니다. attemptsMade reset은 자동 재시도 예산을 다시 주고, attemptsStarted reset은 시작 횟수를 0으로 되돌립니다.

Job.retry()failedReason, finishedOn, processedOn, returnvalue를 지우므로 snapshot이 먼저입니다. 개수·기간 기반 자동 제거는 새 작업이 완료 또는 실패할 때 lazy하게 실행되며, 제거된 job은 상태 조회와 수동 retry가 불가능합니다.


자동 재시도 관찰하기

첫 두 번은 실패하고 세 번째에 성공하는 작업을 제출합니다. requestedBytenantId는 access token으로 확인한 주체에서 채워집니다.

curl -i -X POST http://localhost:3000/reports \
  -H "Content-Type: application/json" \
  -H "Authorization: Bearer <access-token>" \
  -d '{"reportId":"retry-demo","rows":1000,"failUntilAttempt":2}'

응답의 jobId 또는 statusUrl을 사용해 상태를 반복 조회합니다. 자동 증가값을 추측해 /jobs/2처럼 고정하지 않습니다.

JOB_ID='POST 응답의 jobId'
curl -H "Authorization: Bearer <access-token>" \
  "http://localhost:3000/reports/jobs/$JOB_ID"

조회 시점에 따라 delayed, waiting, active 중 일부를 건너뛰어 관찰할 수 있으며, 최종 성공은 state === 'completed'result로 확인합니다.

다음도 확인합니다.

  1. 같은 인증 주체가 같은 reportIdrows를 첫 작업이 끝나기 전에 다시 제출하면 Simple deduplication이 뒤 추가를 억제합니다.
  2. rows처럼 결과를 바꾸는 입력이 달라지면 operation key도 달라져 별도 작업이 됩니다.
  3. 같은 결과가 이미 저장된 뒤 새 작업을 실행하면 멱등 저장소가 기존 결과를 반환하고 reusedtrue가 됩니다.
  4. failUntilAttempt5로 보내면 최대 네 번을 모두 사용한 뒤 failed가 됩니다.

실패 작업을 수동 복구하기

자동 재시도를 모두 사용한 작업은 원인을 고친 뒤 권한이 있는 운영자가 다시 실행할 수 있어야 합니다.

수동 복구는 다음 순서를 지킵니다.

  1. 인증과 retry 권한, 작업 소유 범위를 확인합니다. 이 예제는 tenantIdrequestedBy가 모두 일치해야 합니다. 권한이 없으면 존재 여부를 숨기고 404로 처리합니다.
  2. 같은 operation key의 결과나 새 미완료 작업이 없는지 확인하고, 운영 저장소의 lock 또는 원자적 상태 전이로 복구 경쟁을 막습니다.
  3. payload의 비밀이 아닌 식별자, failedReason, stacktrace, attemptsMade, attemptsStarted, 실행자와 이유를 외부 감사 저장소에 먼저 기록합니다.
  4. 장애 원인이 해결된 뒤에만 Job.retry()를 호출하고 새 상태를 상태 API로 관찰합니다.

BullMQ v5.77.6Job.retry()는 작업을 다시 대기로 옮기면서 failedReason, finishedOn, processedOn, returnvalue를 지웁니다. 따라서 evidence snapshot은 호출 에 남겨야 합니다.

재시도 이유를 검증할 DTO를 추가합니다.

src/reports/dto/retry-report-job.dto.ts
import { IsString, MinLength } from 'class-validator';

export class RetryReportJobDto {
  @IsString()
  @MinLength(10)
  reason: string;
}

아래 recoveryAudit.append()는 데이터베이스나 변경 불가능한 감사 저장소에 append하는 애플리케이션 서비스이고, artifacts는 앞에서 만든 결과 저장소라고 가정합니다. 기존 Nest import에는 BadRequestException, ConflictException, NotFoundException을 합칩니다.

예제의 조회 조건은 결정을 보여 주기 위한 것입니다. 운영 환경에서는 producer와 수동 retry가 함께 참여하는 operation-key lock 또는 원자적 복구 상태 전이 안에서 기존 결과 확인 → 새 미완료 작업 확인 → evidence 기록 → retry를 수행해야 합니다.

src/reports/reports.service.ts (retry 추가)
async retry(
  jobId: string,
  reason: 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 !== 'failed') {
    throw new BadRequestException(
      'only failed jobs can be retried: ' + state,
    );
  }

  const operationKey = reportOperationKey(job.data);
  const liveJobId =
    await this.reportsQueue.getDeduplicationJobId(operationKey);

  if (liveJobId && liveJobId !== job.id) {
    throw new ConflictException(
      'another report job is already in progress',
    );
  }

  if (this.artifacts.find(operationKey)) {
    throw new ConflictException('report result already exists');
  }

  await this.recoveryAudit.append({
    operationKey,
    jobId: job.id,
    tenantId: job.data.tenantId,
    requestedBy: job.data.requestedBy,
    reportId: job.data.reportId,
    rows: job.data.rows,
    failedReason: job.failedReason,
    stacktrace: job.stacktrace ?? [],
    attemptsMade: job.attemptsMade,
    attemptsStarted: job.attemptsStarted,
    retriedBy: principal.userId,
    reason,
    recordedAt: new Date().toISOString(),
  });

  await job.retry('failed', {
    resetAttemptsMade: true,
    resetAttemptsStarted: true,
  });

  return {
    jobId: job.id,
    accepted: true,
    statusUrl: '/reports/jobs/' + job.id,
  };
}

resetAttemptsMade: true는 자동 재시도 판단에 쓰는 attemptsMade를 0으로 되돌려 작업에 attempts: 4의 전체 실행 예산을 다시 줍니다. resetAttemptsStarted: true는 작업이 active로 이동한 횟수인 attemptsStarted를 별도로 0으로 되돌립니다. 두 counter의 reset은 단순 정리가 아니라 감사 기록과 실행 예산을 바꾸는 운영 결정입니다.

소진된 attemptsMade를 reset하지 않으면 재실행에서 다시 오류가 났을 때 추가 자동 재시도 없이 최종 실패합니다. 반대로 원인을 고치지 않고 reset하면 같은 장애를 최대 네 번 더 실행합니다.

Controller에서는 인증 Guard에 더해 retry 전용 권한 Guard를 적용합니다.

src/reports/reports.controller.ts (retry 추가)
@UseGuards(ReportRetryGuard)
@Post('jobs/:jobId/retry')
@HttpCode(HttpStatus.ACCEPTED)
retry(
  @Param('jobId') jobId: string,
  @Body() dto: RetryReportJobDto,
  @Req() request: AuthenticatedRequest,
) {
  return this.reportsService.retry(
    jobId,
    dto.reason,
    request.user,
  );
}

Job.retry()가 반환된 직후 Worker가 작업을 가져갈 수 있으므로 응답은 waiting을 약속하지 않습니다. 202 Accepted와 같은 jobId, statusUrl만 반환하고 실제 상태는 다시 조회합니다.

실패 기록을 너무 빨리 지우지 않는다

removeOnFail: true로 모든 실패를 완료 직후 삭제하면 getJob()과 수동 retry가 불가능하고 원인과 payload도 조사할 수 없습니다.

개수나 기간을 지정한 자동 제거는 새 작업이 completed 또는 failed가 될 때 지연 실행됩니다. 따라서 removeOnComplete: 100removeOnFail: 100은 정확한 보존 시간이 아니라 Redis 사용량을 제한하는 lazy 정책입니다.

제거된 작업은 상태·결과·오류를 더는 조회하거나 재시도할 수 없습니다. 장기 감사와 복구 evidence는 보존 기간을 정한 별도 저장소로 내보냅니다.

정책목적반드시 함께 정할 경계
attempts일시 장애 복구영구 실패 분류와 최대 실행 비용
backoff재시도 부하 분산외부 시스템 복구 시간과 jitter
deduplication.id미완료 작업의 뒤 추가 억제tenant·사용자·작업 버전·결과 입력 범위
operation key부수 효과 중복 방지원자적 unique·conditional write
완료·실패 보관상태 조회와 단기 조사lazy 제거, Redis 한도, 외부 감사 보존
수동 retry원인 수정 뒤 복구권한·tenant·경쟁 방지·evidence·전체 예산

복구 가능한 시스템은 단순히 많이 재시도하는 시스템이 아닙니다.

실패의 성격을 구분하고, 다시 실행해도 결과가 망가지지 않게 만들며, 자동 복구가 끝난 뒤의 권한·감사·보존 절차까지 갖춘 시스템입니다.

다음 절에서는 지연 작업과 Job Scheduler를 추가하고, Worker 동시성·이벤트·메트릭·종료 절차를 운영 관점에서 연결합니다.