Asset State Store API 상세 — 워터마크 패턴과 Watcher Trigger
Apache Airflow Official Documentation (in-repo snapshot) — Apache Software Foundation core-concepts/asset-state-store.rst 전체 (Airflow 3.3에 추가)
이 모듈을 다 읽으면
- asset_state_store가 언제 context에 채워지는지 조건을 설명할 수 있다
- 단일 inlet/outlet일 때의 단축 문법과 다중일 때의 subscript 문법 차이를 안다
- watermark 패턴을 이용한 증분 로드 태스크를 구현 아이디어 수준에서 설명할 수 있다
- BaseEventTrigger 안에서 asset_state_store에 접근하는 방식을 설명할 수 있다
Asset state store는 특정 Dag run과 무관하게 자산(asset) 자체에 스코프된 영구 key/value 저장소로, run을 가로지르는 워터마크나 증분 로드 커서에 쓰인다. 이 모듈은 접근 조건과 문법, get/set/delete/clear API, 워터마크 패턴, 그리고 BaseEventTrigger 안에서의 접근 방식과 생명주기를 다룬다.
asset_state_store가 채워지는 조건
context["asset_state_store"]는 구체적인(concrete) Asset inlet·outlet에 대해서만 채워진다. 태스크는 asset_state_store에 어떤 항목이라도 들어있으려면 최소 하나의 구체적인 inlet 또는 outlet을 선언해야 한다. inlet도 outlet도 없는 태스크가 context["asset_state_store"]에 접근하면 런타임에 KeyError가 발생한다. 이는 @asset 패턴(자산을 outlet으로 암묵적으로 선언)과 @task 패턴(inlets나 outlets를 명시해야 함) 모두에 적용된다.
핵심 포인트
- asset_state_store를 쓰려면 태스크가 최소 하나의 구체적 inlet/outlet Asset을 선언해야 하며, 없으면 KeyError가 발생한다
- @asset 데코레이터는 자산을 outlet으로 암묵적으로 선언하지만, @task는 inlets/outlets를 명시해야 한다
접근 문법: subscript vs 단일 inlet 단축형
일반적인 형태는 context["asset_state_store"][my_asset].get/set/delete/clear처럼, 자산 객체로 subscript해서 그 자산의 state store를 얻는 것이다. 태스크에 구체적인 inlet 또는 outlet이 정확히 하나뿐이라면, subscript 없이 context["asset_state_store"] 자체에서 바로 get/set/delete/clear를 호출하는 단축형을 쓸 수 있다. 태스크에 구체적 inlet/outlet이 둘 이상이면 이 단축형 호출은 ValueError를 일으키므로, 여러 inlet이 있는 태스크에서는 반드시 subscript 형태(context["asset_state_store"][my_asset])를 써야 한다.
핵심 포인트
- inlet/outlet이 정확히 하나뿐이면 subscript 없이 context['asset_state_store'].get(...) 단축형을 쓸 수 있다
- inlet/outlet이 둘 이상이면 단축형 호출 시 ValueError가 발생하므로 반드시 [my_asset]으로 특정해야 한다
API 레퍼런스와 watermark 패턴
get(key, default)는 task state store와 동일한 의미로 저장된 JSON 값 또는 default를 반환한다. set(key, value)는 task state store와 달리 retention 파라미터가 없다 — 값은 명시적으로 삭제되거나 자산이 비활성화될 때까지 남는다. value는 None을 제외한 JSON 호환 타입이어야 한다. delete(key)는 키가 없으면 아무 일도 하지 않는다. clear()는 그 자산의 모든 키를 삭제한다.
Asset state store의 대표 사용 사례는 워터마크 패턴이다: 증분 로드 태스크가 run마다 워터마크를 전진시키며, 그 워터마크는 run이 아니라 자산 자체에 저장되므로 그 자산을 읽거나 쓰는 어떤 태스크든 접근할 수 있다 — BaseEventTrigger를 이용한 자산 '워칭(watching)'을 만들 때 특히 유용하다.
@task(inlets=[orders], outlets=[orders])
def load_new_orders(**context):
asset_state_store = context["asset_state_store"]
watermark = asset_state_store.get("watermark", default="1970-01-01T00:00:00Z")
rows = fetch_orders_since(watermark)
if not rows:
return
upload_to_warehouse(rows)
new_watermark = max(r["created_at"] for r in rows)
asset_state_store.set("watermark", new_watermark)
run마다 이전 run이 남긴 워터마크를 읽어 그 이후의 새 데이터만 가져오고, 워터마크를 갱신한다. Asset state store가 run을 가로질러 지속되므로, 다음 run은 재시도·수동 재실행·스케줄러 재시작을 겪었더라도 정확히 이전 run이 멈춘 지점에서 시작한다.
핵심 포인트
- Asset state store의 set()에는 retention 파라미터가 없다 — 값은 명시적으로 삭제되거나 자산이 비활성화될 때까지 남는다
- 워터마크 패턴은 이전 run이 남긴 워터마크를 읽어 그 이후 데이터만 처리하고 워터마크를 갱신하는 증분 로드의 정석이다
- 워터마크는 재시도, 수동 재실행, 스케줄러 재시작을 가로질러도 마지막 지점에서 이어진다
Watcher Trigger 안에서의 접근과 생명주기
BaseEventTrigger 서브클래스(워처 트리거)는 run() 메서드 안에서 self.asset_state_store로 asset state store를 직접 읽고 쓸 수 있다. Triggerer가 run()이 호출되기 전에 이를 주입해주며, __init__이나 serialize()에서는 접근할 수 없고 오직 run() 안에서만 접근 가능하다. 태스크 기반 접근(자산이 inlet/outlet 선언으로 식별됨)과 달리, 워처 트리거의 접근자는 자동으로 감시 중인 자산에 바인딩되어 있어 subscript가 필요 없다. 트리거가 쓴 값은 그 자산을 inlet이나 outlet으로 선언한 어떤 태스크에서도 보이며, 그 반대도 마찬가지다. 참고로 self.asset_state_store는 BaseEventTrigger 서브클래스에서만 쓸 수 있다 — 태스크 디퍼럴에 쓰이는 일반 BaseTrigger 서브클래스에는 이 접근자가 없다.
생명주기와 가비지 컬렉션: Asset store 행은 무기한 보존되며, task state store에 적용되는 [state_store] default_retention_days 시간 기반 만료 대상이 아니다. 유일한 자동 정리는 orphan sweep이다 — 자산이 비활성화되면(asset_active 레코드가 없어지면) 다음 가비지 컬렉션 패스에서 그 자산의 저장소 행이 제거된다. 그 스윕이 실행되기 전까지는 오래된 행이 DB에 남아있을 수 있지만 더 이상 쓸 수는 없으며, Execution API 리졸버는 활성 자산만 필터링해 보여준다. 명시적으로 제거하려면 태스크 안에서 clear()를 호출하거나 REST API를 쓴다. [state_store] clear_on_success는 asset state store를 지우지 않는다 — asset store는 설계상 run을 가로지르므로, 태스크 레벨 자동 정리가 다음 run이 의존하는 정보를 파괴하면 안 되기 때문이다. 더 이상 필요 없어지면 항상 명시적으로 지워야 한다.
핵심 포인트
- self.asset_state_store는 BaseEventTrigger의 run() 안에서만 쓸 수 있고, __init__/serialize()에서는 접근할 수 없다
- 일반 BaseTrigger(태스크 디퍼럴용)는 asset_state_store에 접근할 수 없다 — BaseEventTrigger 전용이다
- Asset state store는 default_retention_days의 대상이 아니며, 자산이 비활성화되어야 orphan sweep으로 정리된다