Scheduling · Backend

스케줄링과
Task Outbox 패턴

비즈니스 로직을 직접 돌리는 Cron job은 서비스 인스턴스가 하나일 때는 잘 돈다. 문제는 인스턴스가 둘이 되어 같은 작업을 동시에 실행하는 순간부터다. 스케줄링을 제대로 하려면 결국 누가 무엇을 어떤 순서로 해도 되는지를 정해야 한다.

주기 작업과 배치 처리에는 요구사항이 세 가지 따라붙는다. 프로토타입 단계에서는 건너뛰기 쉽지만, 나중에 붙이려면 비싸다. 첫째, Scheduler는 비즈니스 로직이 있는 Application 계층에 두지 않고 Infrastructure 계층에 둔다. 둘째, Task handler는 멱등(idempotent)해야 한다. 메시지 큐는 at-least-once 전달이라 같은 Task가 두 번 돌 수 있다. 셋째, 메시지 큐를 쓴다면 Dead Letter Queue는 나중에 덧붙일 것이 아니고 처음부터 기본으로 둔다. DLQ가 무한 재시도를 끊고, poison message가 뒤의 메시지를 전부 막기 전에 따로 빼 둔다.

Scheduler는 enqueue만 한다

Scheduler는 비즈니스 로직을 직접 실행하지 않는다. 하는 일은 큐에 Task를 넣는 것뿐이다. 실제 작업은 나중에 Task Consumer가 메시지를 받아 Command Service를 호출할 때 일어난다.

[Scheduler] --(enqueue)--> [task_outbox] --(Relay)--> [message queue] --(Consumer)--> [TaskController] --(calls)--> [CommandService]

이렇게 한 단계를 거치면 네 가지를 한꺼번에 얻는다. 먼저 인스턴스가 여러 개여도 안전하다. 여러 인스턴스가 같은 Cron을 동시에 실행해도 FIFO 큐의 중복 제거(deduplication) 덕분에 하나만 처리된다. 재시도도 덤으로 따라온다. Consumer가 실패하면 visibility timeout이 지난 뒤 메시지가 자동으로 다시 전달되고, 최대 수신 횟수를 넘기면 DLQ로 넘어간다.

Backpressure도 생긴다. 작업이 몰리면 큐에 쌓였다가 Consumer의 처리 속도대로 빠져나가니, 하류(downstream) 시스템이 감당 못 할 만큼 밀려들지 않는다. 마지막으로 관측하기 쉽다. 비즈니스 로직에 계측 코드를 넣지 않아도 큐 지표(메시지 수, 처리 지연, DLQ 수)만으로 배치 상태를 알 수 있다.

같은 백엔드 설계를 5개 언어로 나란히 구현해 둔 내 예제 프로젝트가 있다. 그 NestJS 구현의 이자 지급 scheduler를 보면, enqueue 하나만 하고 다른 일은 하지 않는다.

@Injectable()
export class AccountInterestScheduler {
  private readonly logger = new Logger(AccountInterestScheduler.name)

  constructor(private readonly taskQueue: TaskQueue) {}

  @Cron(CronExpression.EVERY_DAY_AT_MIDNIGHT)
  public async enqueueDailyInterest(): Promise<void> {
    const now = new Date()
    const today = new Date(Date.UTC(now.getUTCFullYear(), now.getUTCMonth(), now.getUTCDate()))
    const dateStamp = today.toISOString().slice(0, 10)
    const dedupId = `account.apply-daily-interest-${dateStamp}`

    try {
      await this.taskQueue.enqueue(
        'account.apply-daily-interest',
        { today: today.toISOString() },
        { groupId: 'account.interest', deduplicationId: dedupId }
      )
      this.logger.log({ message: 'Daily interest Task enqueued', dedup_id: dedupId })
    } catch (error) {
      // @nestjs/schedule silently swallows exceptions from Cron handlers, so log explicitly.
      this.logger.error({ message: 'Failed to enqueue daily interest Task', dedup_id: dedupId, error })
    }
  }
}

여러 인스턴스에서 돌아도 안전한 건 날짜를 박은 dedupId 덕분이다. 인스턴스 3개가 같은 FIFO dedup window 안에서 이 handler를 실행해도, 세 시도가 모두 같은 dedupId를 들고 가므로 큐에는 하나만 들어간다.

enqueue 호출을 try-catch로 감싼 이유는 주석에 적힌 그대로다. 여기서 쓰는 스케줄링 라이브러리는 Cron handler 안에서 던진 예외를 아무 말 없이 삼킨다. catch해서 로그를 남기지 않으면 enqueue가 실패해도 흔적 하나 남지 않는다.

Cron 데코레이터가 있을 때와 없을 때

Spring Boot도 같은 cron 표현식을 쓴다. NestJS 데코레이터 대신 표준 라이브러리 애노테이션을 붙인다는 점만 다르다.

@Component
@RequiredArgsConstructor
public class InterestPaymentScheduler {
    private static final String TASK_TYPE = "account.pay-interest";
    private static final String GROUP_ID = "account.interest";
    private final TaskOutboxWriter taskOutboxWriter;

    @Scheduled(cron = "0 0 3 * * *") // Every day at 3 AM
    public void enqueueDailyInterestPayment() {
        try {
            LocalDate today = LocalDate.now();
            String dedupId = TASK_TYPE + "-" + today;
            taskOutboxWriter.enqueue(TASK_TYPE, new Payload(today), GROUP_ID, dedupId);
        } catch (Exception e) {
            log.error("Failed to enqueue the interest-payment Task", e);
        }
    }
}

Go에는 메서드에 붙일 스케줄링 라이브러리가 아예 없다. 그래서 평범한 goroutine 하나가 ticker 루프를 직접 돌린다. 이 루프도 프로세스 안의 다른 루프들과 같은 shutdown context를 지켜본다.

func (s *InterestScheduler) Run(ctx context.Context) {
	ticker := time.NewTicker(24 * time.Hour)
	defer ticker.Stop()

	for {
		select {
		case <-ctx.Done():
			return
		case <-ticker.C:
			if err := s.EnqueueDailyInterest(ctx, time.Now().UTC()); err != nil {
				// Many scheduling libraries silently swallow Cron exceptions — this one
				// always logs it explicitly. Retried on the next tick 24 hours later,
				// so it is not re-thrown here.
				slog.ErrorContext(ctx, "interest payment task enqueue failed", "error", err)
			}
		}
	}
}

func (s *InterestScheduler) EnqueueDailyInterest(ctx context.Context, today time.Time) error {
	date := today.Format("2006-01-02")
	dedupID := "account.apply-interest-" + date
	payload := []byte(`{"date":"` + date + `"}`)
	return s.taskQueue.Enqueue(ctx, "account.apply-interest", payload, dedupID)
}

세 구현이 프레임워크에서 받는 도움은 데코레이터, 애노테이션, 손으로 짠 ticker로 제각각이다. 그런데 속을 보면 셋 다 같은 모양이다. enqueue만 하고, 스택 어딘가가 예외를 삼키기 쉬우니 실패는 직접 로그로 남긴다. 인스턴스끼리 조율하려 들지 않고, 날짜 기반 dedup ID가 다중 인스턴스 상황을 흡수하게 둔다.

Enqueue는 DB 변경과 원자적이어야 한다

Command Service 안에서 메시지 큐의 SendMessage를 바로 호출하면, 신뢰성 있는 이벤트 기반 설계에서 늘 나오는 dual-write 문제가 그대로 생긴다. DB는 커밋됐는데 메시지 전송이 실패하거나, 메시지는 나갔는데 DB가 롤백된다. 그러면 아무도 지켜보지 않는 불일치가 남는다.

해법도 Domain Event 때와 같은 Outbox 패턴이다. DB 변경과 같은 트랜잭션 안에서 task_outbox 테이블에 쓰고, 별도의 Relay가 그 테이블을 폴링해 트랜잭션이 커밋된 뒤에만 발행한다.

// An Application Service — the DB change and enqueuing the Task happen in the same transaction
await transactionManager.run(async () => {
  await orderRepository.saveOrder(order)
  await taskQueue.enqueue(
    'order.archive',
    { orderId: order.orderId },
    { groupId: order.orderId, deduplicationId: `order.archive-${order.orderId}` }
  )
})

Cron tick에서 실행되는 Scheduler처럼 트랜잭션 컨텍스트가 아예 없는 곳에서도 같은 경로를 쓴다. row 하나를 insert하는 일이니 그 자체로 원자적이다. enqueue하는 곳마다 경로를 하나로 맞춰 두면 머릿속 모델도 단순해진다. 무엇이 계기였든 enqueue는 언제나 outbox 테이블에 쓰는 일이고, 큐 클라이언트를 직접 부르는 일은 없다.

Task Controller는 Handler가 아니라 Interface 계층의 Adapter다

HTTP Controller가 HTTP 요청을 받아 Application Service에 넘기듯, Task Controller는 메시지 큐의 메시지를 받아 Command Service를 호출한다. 자기만의 조건 분기나 비즈니스 규칙은 없다. HTTP Controller와 다른 점은 에러를 잡아서 변환하지 않는다는 것이다. 예외를 그대로 다시 던진다. 그 예외를 재시도로 돌릴지 DLQ로 보낼지는 Consumer가 정하기 때문이다.

class OrderTaskController {
  constructor(private readonly orderCommandService: OrderCommandService) {}

  async archive(payload: ArchiveOrderCommand): Promise<void> {
    await this.orderCommandService.archiveOrder(payload)  // the exception propagates as-is
  }
}

멱등성(Idempotency)의 세 단계

전달이 at-least-once이니 Task handler는 몇 번을 돌든 같은 결과를 내야 한다. Level 1은 로직 자체가 멱등한 경우다. 이미 만료된 주문을 archive하는 작업이 그렇다. 이미 archive된 주문을 다시 처리해도 아무 일도 일어나지 않는다(no-op). Level 2는 원장(ledger)을 쓴다. side effect가 있는 handler가 어떤 ID를 처리했는지 기록해 두고, 중복이 오면 건너뛴다. Level 3는 강한 원자성이 필요한 경우다. handler 로직과 ledger 기록을 한 트랜잭션으로 묶어서, 중간에 실패해도 한쪽만 기록되는 일이 없게 한다.

실제 인프라에서 나온 버그들

5개 언어 구현 모두에 이 기능을 넣으면서, 설계 리뷰만으로는 안 나왔을 버그를 몇 개 만났다. 모두 실제 인프라에서 테스트를 동시에 돌려 봐야 드러나는 것들이었다. 한 구현에서는 SQS task-queue 설정에 필수 config 필드가 새로 생겼다. 그런데 end-to-end 테스트 클래스 6개 중 5개가 이 값을 넣지 않아서, 그 테스트들에서만 앱 부팅이 깨졌다.

다른 구현의 SQS FIFO 충돌은 훨씬 알아채기 어려웠다. 한 번의 테스트 실행 안에서 여러 테스트 메서드가 같은 월간 scheduler를 불렀는데, 모두 같은 날짜 기반 dedupId를 만들었다. dedup window는 분 단위라서 첫 호출만 큐에 닿았고, 나머지는 중복으로 걸러져 에러 없이 사라졌다. 그 바람에 거짓 양성(false-positive)으로 테스트가 통과할 뻔했다. 한 시나리오의 assertion이 우연히 맞아떨어졌는데, 그 assertion이 기대던 시나리오는 한 번도 돌지 않았다.

두 버그의 공통점

인메모리 fake 큐로 돌리는 유닛 테스트는 실제 dedup window도, 실제 필수 config 검사도 겪지 않는다. fake가 둘 다 강제하지 않으니 겪을 수가 없다. 스케줄링 기능은 실제 인프라에서 동시 호출로 적어도 한 번은 돌려 봐야 검증했다고 할 수 있다.

Payload는 작게

SQS는 메시지 하나를 256KB로 제한한다. 넘을 수 없는 상한이니, 장애가 나서야 알게 되기보다 처음부터 설계에 넣어 두는 게 낫다. Payload에는 { orderId: 'o1' } 같은 작은 메타데이터만 담는다. 큰 데이터는 S3에 올리고 메시지에는 storage key만 싣는다.

더 볼 자료

docs/architecture/scheduling.md(예제 프로젝트의 Task Outbox 패턴 전체, MessageGroupId 전략, DLQ 모니터링) · account-interest-scheduler.ts(위에서 본 scheduler를 코드 맥락 그대로)