API 밖에서 도는 worker와 CronJob

수집, 예측과 푸시 작업을 요청 경로에서 분리하고 중복 실행을 막은 방법

한강자리 백엔드는 사용자가 부르는 API와 뒤에서 도는 worker를 따로 띄운다. 코드는 같지만 맡은 일이 다르다.

API는 사용자가 앱을 열면 바로 답해야 한다. 공공데이터 수집, 예측 계산과 푸시 대상 선정은 느리거나 실패할 수 있고 실행 간격도 달라 별도 worker가 맡는다.

초기에는 백엔드 프로세스 하나에 scheduler를 붙이는 방안도 검토했다. 운영 경로에서는 API가 짧은 요청에 바로 답하고 worker는 실패한 작업을 다시 시도하며 실행 기록을 남겨야 했다.

worker는 실패 범위를 작게 가두도록 나눴다. 주차 상태 확인이 잠시 실패해도 API는 마지막 확인값을 내려주고, push outbox가 밀려도 home-summary 응답은 영향을 받지 않아야 한다.

API 밖에서 미리 하는 일

지금 계속 떠 있는 worker는 네 종류다. 이들은 사용자가 기다리는 API 호출 밖에서 값을 먼저 만들어 둔다.

flowchart TB
  Scheduler["워커별 APScheduler"] --> Parking["worker-parking\n상태 폴링"]
  Scheduler --> Outing["worker-outing\n행사 / 공지 / 도시 맥락 / 번역"]
  Scheduler --> Forecast["worker-forecast\n예측 생성 / 첫 화면 워밍업 / 메트릭 집계"]
  Scheduler --> Push["worker-push\n후보 / 선호 인덱스 / outbox / 유지보수"]

  Parking --> PG["Postgres"]
  Parking --> Redis["Redis"]
  Outing --> PG
  Forecast --> PG
  Forecast --> Redis
  Push --> PG
  Push --> Redis
  Push --> APNs["APNs"]

가끔 기준 데이터를 맞추거나 값이 맞는지 확인하는 일은 CronJob 쪽으로 따로 본다. 상시 polling과 기준 데이터 sync를 같은 그림에 넣으면, 무엇이 자주 돌고 무엇이 드물게 도는지 흐려진다.

flowchart LR
  Cron["K8s CronJob"] --> Master["주차 기준 데이터 동기화"]
  Cron --> Facility["시설 동기화"]
  Cron --> Transit["대중교통 데이터셋 동기화"]
  Cron --> Backtest["예측 백테스트"]
  Master --> PG
  Facility --> PG
  Transit --> PG
  Backtest --> PG

헷갈리기 쉬운 점도 있다. worker-parking이라는 이름만 보면 master sync와 status polling이 모두 상시 scheduler에 있는 것처럼 보일 수 있다. 현재 코드에서 상시 scheduler에 등록된 것은 status polling이다. parking master sync는 별도 job entrypoint와 K8s CronJob 쪽 일정으로 분리되어 있다.

실행 간격과 작업의 분리

코드에서는 일을 다음처럼 나눈다.

worker주요 작업
worker-parkingstatus polling
worker-outingevent sync, notice sync, realtime context sync, translation sync
worker-forecastforecast generation, home-summary precompute, metric rollup
worker-pushcandidate build, preference index rebuild, outbox drain, maintenance
K8s CronJob기준 데이터 동기화, 시설 동기화, 데이터셋 동기화, backtest

이렇게 나누면 한쪽 장애가 다른 화면까지 끌고 가지 않는다. 주차 데이터를 가져오는 쪽이 잠시 실패해도 API는 마지막으로 확인한 상태를 내려줄 수 있다. 푸시 발송이 지연되어도 주차 화면이 값을 읽는 쪽과는 떨어져 있다.

운영자는 worker 이름에 따라 첫 확인 대상을 고른다. 주차 수집 쪽이 흔들리면 최신성부터 보고, 알림 작업이 밀리면 대기열과 전송 기록을 본다. 하나의 worker에 모든 일이 들어 있으면 어디부터 봐야 할지 늦어진다.

중복 실행 방지

주차 상태 확인을 단순화하면 다음과 같다.

sequenceDiagram
  autonumber
  participant Scheduler as Scheduler
  participant Job as 상태 폴링 작업
  participant Lock as DB 잠금
  participant Source as 주차 API
  participant PG as Postgres
  participant Redis as 상태 캐시
  participant Fact as 알림 fact

  Scheduler->>Job: 정해진 간격에 실행
  Job->>Lock: 중복 실행 잠금 시도
  alt 잠금 획득
    Job->>Source: 주차 상태 원천 데이터 요청
    Source-->>Job: 원천 응답 행 반환
    Job->>PG: 상태 스냅샷과 수집 기록 저장
    Job->>Redis: 최신 주차 상태 캐시 갱신
    Job->>Fact: 변화가 있으면 알림 후보 기록
  else 중복 실행
    Job-->>Scheduler: 이번 실행 건너뛰기
  end

worker는 잠금과 실행 기록으로 중복 실행을 확인한다. 같은 일이 동시에 두 번 돌면 중복 row, 중복 알림, cache 경합이 생길 수 있다. 그래서 scheduler 설정만 믿지 않았다.

저장소 수준에도 잠금을 두어 같은 작업이 동시에 처리되지 않게 했다.

데이터별 수집 주기

한강자리에는 데이터를 가져오는 곳마다 갱신 주기를 적어 둔 schedule catalog가 있다. 여기에는 출처, 작업 소유 범주, 갱신 간격, stale로 볼 시간 같은 정보가 들어간다.

운영자는 이 catalog에서 출처별 실행 주기와 stale 기준을 확인한다. API의 최신성 표시와 운영 알림도 같은 기준을 사용한다.

30초 단위로 확인하는 값과 하루 한 번 확인하는 값에는 서로 다른 최신성 문구가 필요하다. catalog가 수집 간격과 stale 기준을 한곳에 묶어 이 차이를 유지한다.

flowchart LR
  Catalog["수집 계획표"] --> SchedulerJobs["상시 scheduler 작업"]
  Catalog --> CronJobs["운영 CronJob"]
  SchedulerJobs --> Ingestion["수집 유스케이스"]
  CronJobs --> Ingestion
  Ingestion --> Runs["ingestion_runs\n성공 / row_count / schema_hash"]
  Runs --> Health["출처 상태와 최신성"]
  Health --> API["나들이/home-summary 최신성"]

작업 설계의 공통 원칙

  • API startup hook에 ingestion을 넣지 않는다.
  • 작업은 가능한 idempotent하게 만든다.
  • 데이터를 가져오는 곳마다 success, failure, row count, schema hash를 남긴다.
  • 중복 실행은 scheduler 설정과 저장소 잠금으로 막는다.
  • 푸시 발송은 outbox claim과 delivery attempt를 보고 추적할 수 있게 만든다.
  • 데이터를 가져오는 곳의 실패는 앱 전체 장애가 아니라 해당 영역의 최신성으로 드러낸다.

API와 worker를 나눈 결과

API는 마지막 확인값을 빠르게 돌려주고, worker는 출처별 일정에 따라 새 값을 수집한다. 실행 생명주기를 분리하자 수집 실패와 사용자 요청 지연을 따로 추적할 수 있었다.

Comments

댓글

    이미지 확대