← 학습 카테고리

Learn

Airflow

151개 모듈 · 현재 30번째

Airflow 모듈 30/151 airflow-learn-30

Task State Store API 상세 — get/set/delete/clear와 활용 패턴

Apache Airflow Official Documentation (in-repo snapshot) — Apache Software Foundation core-concepts/task-state-store.rst 전체 (Airflow 3.3에 추가)

이 모듈을 다 읽으면

  • context['task_state_store']의 동기/비동기 메서드를 사용할 수 있다
  • retention과 NEVER_EXPIRE의 차이를 설명할 수 있다
  • 외부 작업 재접속, 페이지네이션 체크포인트, 진행률 노출의 세 가지 활용 패턴을 구현 아이디어 수준에서 설명할 수 있다
  • 매핑된 태스크와 clear_on_success가 저장소 동작에 미치는 영향을 설명할 수 있다

Task state store는 (dag_id, run_id, task_id, map_index)에 스코프된 영구 key/value 저장소로, context['task_state_store']를 통해 접근한다. 이 모듈은 get/set/delete/clear와 async 대응 메서드, retention/NEVER_EXPIRE, 외부 작업 재접속·페이지네이션 체크포인트·진행률 노출 패턴, 그리고 매핑된 태스크와 clear_on_success의 동작을 다룬다.

기본 API: get/set/delete/clear

task_state_store는 (dag_id, run_id, task_id, map_index)로 스코프된 태스크 인스턴스 단위의 영구 key/value 저장소다. 워커 크래시와 같은 Dag run 안에서의 태스크 재시도를 버텨내므로, 외부 작업 ID·태스크 내 체크포인트·진행 상황 메타데이터를 저장하기에 적합하다. context['task_state_store']를 통해 접근하며, 동기 메서드 get/set/delete/clear와, async 태스크용 비동기 대응 메서드 aget/aset/adelete/aclear를 제공한다.

get(key, default)는 저장된 JSON 값을 반환하거나, 키가 없으면 default를 반환한다. set(key, value, *, retention=None)는 값을 쓰거나 덮어쓴다. value는 None을 제외한 JSON 호환 타입(str, int, float, bool, list, dict)이면 무엇이든 가능하다. retention 인자는 키의 만료 시점을 제어한다: timedelta(...)를 주면 쓰기 시점으로부터 그 기간 뒤 만료되며(만료 시각은 값을 API 서버로 보내기 전 워커에서 계산된다), NEVER_EXPIRE 센티널을 주면 전역 [state_store] default_retention_days 설정과 무관하게 절대 만료되지 않고 가비지 컬렉션에서 제외되며, None(기본값)이면 전역 default_retention_days 설정을 따른다. 중요한 점은 retention이 timedelta만 허용한다는 것이다 — 정수 일수를 그냥 넘기면 TypeError가 발생한다. delete(key)는 단일 키를 삭제하며 키가 없으면 아무 일도 하지 않는다. clear()는 이 태스크 인스턴스의 모든 키를 삭제한다.

핵심 포인트

  • set의 retention은 timedelta만 허용하며 정수 일수를 넘기면 TypeError가 발생한다
  • NEVER_EXPIRE는 전역 default_retention_days 설정과 무관하게 가비지 컬렉션에서 제외된다
  • clear()는 해당 태스크 인스턴스의 모든 키를 지운다 — 특정 키만 지우려면 delete(key)를 쓴다

비동기 접근자 (aget/aset/adelete/aclear)

aget, aset, adelete, aclear는 동기 대응 메서드와 같은 인자를 받고 동일하게 동작하지만, 이벤트 루프를 블로킹하는 대신 API 서버와의 왕복을 await한다 — 그래서 코루틴이 다른 동시 작업을 멈추지 않고도 진행 상황을 체크포인트할 수 있다. async 태스크 안에서 동기 get/set/delete/clear를 호출하면 이벤트 루프를 블로킹해, 그 코루틴이 애초에 얻으려던 동시성 이점을 무너뜨린다 — 반드시 a로 시작하는 메서드를 써야 한다.

핵심 포인트

  • 비동기 태스크 안에서 동기 get/set을 호출하면 이벤트 루프가 블로킹되어 코루틴의 동시성 이점이 사라진다 — 반드시 aget/aset을 써야 한다

활용 패턴: 외부 작업 재접속과 페이지네이션 체크포인트

외부 작업 재접속 패턴에서는, 작업을 제출하기 전에 이미 저장된 job_id가 있는지 먼저 확인한다. 없으면(None) 새로 제출하고, 작업이 끝나기 전에 키가 GC되지 않도록 retention=NEVER_EXPIRE로 저장한다. 그런 다음 저장된 ID로 결과를 재연결·대기한다.

job_id = task_state_store.get("job_id")
if job_id is None:
    job_id = spark_client.submit_job(...)
    task_state_store.set("job_id", job_id, retention=NEVER_EXPIRE)
result = spark_client.wait_for_completion(job_id)

재시도 시에는 저장된 job_id를 발견해 중복 작업을 제출하는 대신 재연결한다. BaseOperator 서브클래스라면 이 패턴을 ResumableJobMixin이 캡슐화해준다 — 제출 후 외부 작업 ID를 task state store에 저장해두고, 재시도 시 활성 작업이 있으면 재연결하고 이전 작업이 최종 실패 상태였을 때만 재제출한다.

페이지네이션/배치 처리되는 데이터를 다루는 태스크에서는 마지막으로 완료한 오프셋(예: last_page)을 저장해두면, 재시도가 처음부터가 아니라 중간부터 재개된다. start_page는 저장된 last_page + 1(없으면 1)로 계산해 그 지점부터 순회한다. 같은 패턴을 async 태스크에서는 aget/aset으로 구현해, 페이지네이션되는 외부 API를 순회하는 코루틴이 이벤트 루프를 막지 않고 체크포인트하게 할 수 있다.

핵심 포인트

  • 외부 작업 재접속 패턴은 job_id를 NEVER_EXPIRE로 저장해두고, 있으면 재접속·없으면 새로 제출한다
  • ResumableJobMixin은 이 패턴을 캡슐화해 재시도 시 활성 작업에 재접속하거나 종료 실패 상태였을 때만 재제출한다
  • 페이지네이션 체크포인팅은 last_page 같은 오프셋을 저장해 재시도가 중간부터 이어가게 한다

진행률 노출, 동기/디퍼러블 차이, 매핑된 태스크, 자동 정리

진행률 메타데이터 패턴에서는 task_state_store.set("progress", {"rows_loaded": total, "status": "running"})처럼 row 카운트·상태 문자열·가벼운 JSON을 저장한다. 이 progress 키는 XCom이나 외부 시스템 없이도 REST API와 Airflow UI를 통해 태스크가 실행되는 동안 실시간으로 노출된다.

동기 태스크와 디퍼러블 태스크는 동작이 조금 다르다. 동기 태스크는 워커 프로세스가 크래시하면 태스크 인스턴스가 재시도되는데, 크래시 전에 쓴 state store 데이터는 보존되어 재시도가 이어받을 수 있다(위의 외부 작업 재접속 패턴 참고). 디퍼러블 태스크는 일단 defer되면 Triggerer가 poke 주기 사이의 연속성을 처리하므로, task state store는 일반적인 poke 연속성이 아니라 오퍼레이터가 시작한 clear를 견뎌내야 할 때만 쓴다.

매핑된 태스크(task.expand(...))에서는 각 map index가 자신만의 task state store 네임스페이스를 갖는다. clear()는 현재 index의 저장소만 지운다. 모든 map index에 걸쳐 상태를 지우려면 태스크 그룹이 끝난 뒤 Core API(UI나 CLI)를 써야 한다.

자동 정리는 [state_store] clear_on_success = True로 설정하면, 태스크 인스턴스가 success 상태로 전환될 때 그 인스턴스의 task state store 키가 모두 자동 삭제된다 — 성공 이후의 관측성이 필요 없을 때 저장 공간을 줄이는 데 유용하다. 단, clear_on_success는 task state store에만 적용된다 — asset state store는 태스크가 아니라 자산에 스코프되어 있어 이 설정의 영향을 절대 받지 않으며, run을 가로질러 지속되므로 명시적으로 지워야 한다.

핵심 포인트

  • progress 같은 키는 REST API와 UI에서 실시간으로 조회할 수 있어 XCom 없이 관측성을 제공한다
  • 매핑된 태스크는 map index마다 독립된 네임스페이스를 가지며, clear()는 현재 index만 지운다 — 전체를 지우려면 Core API를 써야 한다
  • clear_on_success=True는 task state store만 지우며 asset state store는 절대 건드리지 않는다