Task State Store
Task State Store
단일 Task Instance에 스코프된 영속 key/value 저장소인 Task store를 설명하는 문서예요. context["task_state_store"]로 접근하는 get·set·delete·clear 네 가지 메서드와, 외부 job 재개·태스크 내 체크포인팅·진행 메타데이터 같은 유스케이스, 그리고 sync/deferrable·mapped 태스크에서의 동작 차이를 살펴볼게요.
출처: 문서
본문
3.3 버전에 추가됨 (Added in version 3.3).
Task store는 단일 task instance(dag_id + run_id + task_id + map_index)에 스코프된 영속 key/value 저장소예요. 워커 크래시와 같은 Dag run 내의 태스크 재시도를 견디므로, 외부 job ID, 태스크 내 체크포인트, 진행 메타데이터를 저장하는 데 적합해요.
워커 크래시보다 오래 살아남고 다음 시도에서도 계속 읽을 수 있기 때문에, task state store는 durable execution의 메커니즘이에요. 이때 태스크는 작업을 반복하거나 중복 외부 job을 제출하는 대신 멈춘 지점에서 이어서 진행해요. durable 파라미터를 광고하는 provider operator들은 여기 설명된 API 위에 구축돼요.
task state store를 통해 영속화된 데이터는 태스크 context를 통해 context["task_state_store"]로 접근하고, get, set, delete, clear 네 가지 메서드를 노출해요.
task state store 접근하기 (Accessing task state store)
@task 데코레이터가 붙은 함수나 BaseOperator.execute() 메서드 안에서 task state store는 context 딕셔너리를 통해 task_state_store 키로 사용할 수 있어요. 여기에서 특정 key-value 쌍의 데이터를 검색·설정·삭제·지우는 데 사용할 수 있어요. 이 예시에서는 job_id를 task state store에서 검색한 다음 업데이트하고, 삭제하기 전에 진행해요. 그 태스크의 모든 데이터는 clear 메서드로 제거돼요.
from airflow.sdk import task
import random
@task
def my_task(**context):
# Retrieve task_state_store from context
task_state_store = context["task_state_store"]
my_value = task_state_store.get("my_key", default="my_default_key")
# Set the new value
new_value = f"It is {random.randint(1, 12 + 1)} o'clock"
task_state_store.set("my_key", new_value)
# Delete the value
task_state_store.delete("my_key")
# Clear all store entries for the task
task_state_store.clear()
참조 (Reference)
get(key, default)
저장된 JSON 값을 반환하거나, 키가 존재하지 않으면 default 값을 반환해요.
value = task_state_store.get(
"job_id", default="123456789"
) # returns the value associated with `job_id` or the default value
set(key, value, *, retention=None)
지정된 키에 값을 쓰거나 덮어써요. 참고로 value는 None을 제외한 모든 JSON 호환 타입일 수 있어요. 여기에는 다음이 포함돼요:
strintfloatboollistdict
선택적 retention 인자는 키가 언제 만료되는지 제어해요:
timedelta(...): 쓰기 시점부터 주어진 기간이 지나면 만료돼요(예:timedelta(hours=6)). 만료 타임스탬프는 값이 API 서버로 보내지기 전에 worker에서 계산돼요.NEVER_EXPIRE: 키가 절대 만료되지 않고, 전역[state_store] default_retention_days설정과 무관하게 가비지 컬렉션에서 건너뛰어져요.None(기본값): 전역[state_store] default_retention_days구성으로 폴백해요.
중요 (Important)
retention은 timedelta만 받고, 평범한 정수 일(day) 수는 받지 않아요. 정수를 전달하면TypeError가 발생해요.
# correct
task_state_store.set("key", "val", retention=timedelta(days=7))
# wrong — raises TypeError
task_state_store.set("key", "val", retention=7)
NEVER_EXPIRE 센티널 (NEVER_EXPIRE sentinel)
airflow.sdk에서 NEVER_EXPIRE를 import해요:
from airflow.sdk import NEVER_EXPIRE
task_state_store.set("job_id", job_id, retention=NEVER_EXPIRE)
delete(key)
단일 키를 삭제해요. 키가 존재하지 않으면 no-op이에요.
task_state_store.delete("job_id")
clear()
이 task instance의 모든 task state store 키를 삭제해요.
task_state_store.clear()
예시 유스케이스 (Some Example Use Cases)
외부 job 재개 (External job resumption)
오래 실행되는 외부 job의 흔한 패턴: 제출하기 전에 job ID가 이미 저장되어 있는지 확인하고, NEVER_EXPIRE를 사용해 키가 기본 보존 기간보다 오래 살아남게 해요.
from datetime import timedelta
from airflow.sdk import DAG, task
from airflow.sdk import NEVER_EXPIRE
with DAG("spark_job_dag", schedule=None):
@task
def run_spark_job(**context):
task_state_store = context["task_state_store"]
# Check for an already-submitted job from a previous attempt.
job_id = task_state_store.get("job_id")
if job_id is None:
job_id = spark_client.submit_job(...)
# Store with NEVER_EXPIRE so the key is not garbage-collected before the job finishes
task_state_store.set("job_id", job_id, retention=NEVER_EXPIRE)
# Reattach to the job and wait for completion.
result = spark_client.wait_for_completion(job_id)
return result
재시도 시 태스크는 저장된 job_id를 찾아 중복 job을 제출하는 대신 다시 연결해요. 이런 로직의 또 다른 예시는 example_task_state_store.py에서 찾을 수 있어요.
BaseOperator 서브클래스의 경우 ResumableJobMixin이 이 패턴을 캡슐화해요. 제출 후 외부 job ID를 task state store에 영속화하고, 재시도 시 활성 job에 다시 연결하거나 이전 job이 터미널 실패 상태에 도달하면 재제출해요.
태스크 내 체크포인팅 (Intra-task checkpointing)
페이지네이션 또는 배치 데이터를 처리하는 태스크의 경우 마지막으로 완료된 offset을 저장해, 재시도가 처음부터 다시 시작하는 대신 흐름 중간에서 재개할 수 있게 해요.
from airflow.sdk import DAG, task
with DAG("paginated_ingest", schedule="@daily"):
@task
def ingest_pages(**context):
# Retrieve the task_state_store
task_state_store = context["task_state_store"]
raw = task_state_store.get("last_page")
start_page = raw + 1 if raw is not None else 1
for page in range(start_page, total_pages + 1):
fetch_and_load(page)
task_state_store.set("last_page", page) # Update the task_state_store for reuse later
재시도 시 태스크는 last_page를 읽고 이미 처리된 페이지를 건너뜁니다.
진행 메타데이터 (Progress metadata)
Task store는 XCom이나 외부 시스템 없이 관측 가능성을 위한 진행 중 메트릭(행 개수, 상태 문자열, 가벼운 JSON 페이로드)을 노출할 수 있어요.
from airflow.sdk import DAG, task
with DAG("row_ingest", schedule="@hourly"):
@task
def ingest_rows(**context):
task_state_store = context["task_state_store"]
total = 0
for batch in get_batches():
load(batch)
total += len(batch)
task_state_store.set(
"progress",
{"rows_loaded": total, "status": "running"},
)
task_state_store.set(
"progress",
{"rows_loaded": total, "status": "done"},
)
progress 키는 태스크가 실행되는 동안 REST API와 Airflow UI를 통해 볼 수 있어요.
Sync vs. deferrable 태스크
Task store는 태스크가 동기적으로 실행되느냐 deferral 메커니즘을 사용하느냐에 따라 약간 다르게 동작해요.
동기 태스크 (Synchronous tasks)
worker 프로세스가 크래시하면 task instance가 재시도돼요. 크래시 전에 작성된 task store 데이터는 보존되므로, 재시도는 이전 시도가 멈춘 지점부터 이어갈 수 있어요(위 '외부 job 재개' 패턴 참고).
Deferrable 태스크
태스크가 defer하면 Triggerer가 poke 사이클 전반의 연속성을 처리해요. operator가 시작한 clear에서 살아남아야 할 때만 deferrable 태스크에서 task state store를 사용하고, 일반 poke 연속성에는 사용하지 마세요.
Mapped 태스크
태스크가 동적으로 매핑되면(task.expand(...)), 각 map index는 자체 task state store 네임스페이스를 가져요. clear()는 현재 index의 store만 지워요.
태스크의 모든 map index에 걸친 상태를 제거하려면 태스크 그룹이 끝난 후에 Core API(예: UI 또는 CLI 경유)를 사용해요.
# Inside a mapped task — clears only this index
task_state_store.clear()
자동 정리 (clear_on_success)
[state_store] clear_on_success = True일 때, task instance의 모든 task state store 키는 태스크가 success 상태로 이동할 때 자동으로 삭제돼요. 이는 성공 후 관측 가능성이 필요 없을 때 저장 공간을 줄이는 데 유용해요.
참고 (Note)
clear_on_success는 task state store만 지워요. Asset store는 task instance가 아니라 asset에 스코프되며 이 설정의 영향을 결코 받지 않아요. Asset store는 실행 간에 영속되며 명시적으로 지워야 해요.
전체 구성 세부 사항은 Task and Asset State Store Configuration을 참고하세요.