Asset 상태 저장소
Asset 상태 저장소 (Asset State Store)
이 페이지는 버전 3.3에 추가된 Asset 상태 저장소(asset state store)를 다뤄요. asset store는 특정 DAG run과 무관하게 asset에 범위가 지정된 영구 key/value 저장소예요. 수위 표시(watermark), 증분 로드 커서, asset별 설정 같은 여러 run에 걸친 메타데이터를 담는 자연스러운 장소예요. context["asset_state_store"]로 접근하며 전체 API와 예제를 설명해요.
출처: 문서
본문
버전 3.3에 추가됨.
Asset store는 특정 DAG run과 무관하게 asset에 범위가 지정된 영구 key/value 저장소예요. 단일 task instance에 묶여 있는 task state store와 달리, asset state store는 run을 넘어 영속되며 논리적으로 asset 자체가 소유해요. 수위 표시, 증분 로드 커서, asset별 설정 같은 여러 run에 걸친 메타데이터의 자연스러운 장소예요.
Asset store는 Task 컨텍스트를 통해 context["asset_state_store"]로 접근해요.
asset_state_store는 언제 사용 가능한가요?
Task 안에서 asset state store를 사용할 때, context["asset_state_store"]는 구체적(concrete) Asset inlet·outlet에 대해 채워져요. asset_state_store가 어떤 항목을 포함하려면 Task가 적어도 하나의 구체적 inlet 또는 outlet을 선언해야 해요.
Task에 inlet·outlet이 없고 context["asset_state_store"]에 접근하면 런타임에 KeyError가 발생해요. asset state store를 사용하려면 Task에 asset inlet·outlet을 하나 이상 선언해요.
이는 @asset 패턴과 @task 패턴 모두에 적용돼요:
from airflow.sdk import Asset, DAG, asset, task
my_asset = Asset("my_data", uri="s3://bucket/my_data")
# @asset Dags implicitly declare the asset as an outlet
@asset
def my_asset_dag(**context):
context["asset_state_store"].set("watermark", "2024-06-01")
# @task within a @dag requires explicit inlets or outlets
with DAG("example", schedule=None):
@task(outlets=[my_asset])
def my_task(**context):
context["asset_state_store"].set("watermark", "2024-06-01")
context로 asset state store 접근하기
Asset이 inlet이나 outlet에 포함되면 context["asset_state_store"]로 사용 가능해져요. 그런 다음 자산 객체로 context["asset_state_store"]를 subscripting해 그 asset의 asset state store를 검색할 수 있어요.
from airflow.sdk import Asset, DAG, task
my_asset = Asset("my_data", uri="s3://bucket/my_data")
with DAG("example_asset_store", schedule=None):
@task(inlets=[my_asset], outlets=[my_asset])
def process(**context):
asset_state_store = context["asset_state_store"][my_asset]
watermark = asset_state_store.get("watermark")
asset_state_store.set("watermark", "2024-06-01")
실제 DAG에서 asset state store의 동작을 보려면 example_asset_state_store.py의 DAG를 확인해요.
단일 inlet 축약(shorthand)
정확히 하나의 구체적 inlet 또는 outlet이 있는 Task의 경우, subscripting 없이 context["asset_state_store"]에서 get, set, delete, clear를 직접 호출할 수 있어요.
@task(inlets=[my_asset], outlets=[my_asset])
def process_single(**context):
asset_state_store = context["asset_state_store"]
watermark = asset_state_store.get("watermark")
asset_state_store.set("watermark", "2024-06-01")
Task에 구체적 inlet·outlet이 두 개 이상이면 축약을 호출하면 ValueError가 발생해요. Task에 여러 inlet이 있을 때마다 subscript 형식(context["asset_state_store"][my_asset])을 사용해요.
API 참조
다음 메서드는 per-asset 접근자(context["asset_state_store"][my_asset])와 Task에 정확히 하나의 inlet이 있을 때의 축약(context["asset_state_store"]) 모두에서 사용할 수 있어요. 각각에는 async Task와 watcher trigger 안에서 사용하는 a-접두사 비동기 대응물 — aget, aset, adelete, aclear — 이 있어요.
get(key, default)
저장된 JSON 값을 반환하거나, 키가 존재하지 않으면 default 값을 반환해요.
# Using context
watermark = context["asset_state_store"][my_asset].get("watermark", default="initial_watermark")
set(key, value)
key-value 쌍을 쓰거나 덮어써요. task state store와 달리 asset state store에는 retention 파라미터가 없어요. 값은 명시적으로 삭제되거나 asset이 비활성화될 때까지 영속돼요. task state store처럼 value는 None을 제외한 모든 JSON 호환 타입일 수 있어요. 여기에는 다음이 포함돼요:
strintfloatboollistdict
# Using context
context["asset_state_store"][my_asset].set("watermark", "2024-06-01T00:00:00Z")
delete(key)
단일 키를 삭제해요. 키가 존재하지 않으면 no-op이에요.
# Using context
context["asset_state_store"][my_asset].delete("watermark")
clear()
Asset의 모든 asset state store 키를 삭제해요.
# Using context
context["asset_state_store"][my_asset].clear()
aget, aset, adelete, aclear
get, set, delete, clear의 비동기 대응물. 같은 인자를 받고 동기 형제들과 동일하게 동작해요 — aset에도 retention 파라미터가 없어요 — 하지만 이벤트 루프를 블로킹하는 대신 API 서버로의 왕복을 await하므로, 코루틴이 다른 동시 작업을 막지 않고 asset 상태를 읽고 진행할 수 있어요.
watermark = await context["asset_state_store"][my_asset].aget("watermark", default="initial_watermark")
await context["asset_state_store"][my_asset].aset("watermark", "2024-06-01T00:00:00Z")
await context["asset_state_store"][my_asset].adelete("watermark")
await context["asset_state_store"][my_asset].aclear()
async task 안에서 동기 get/set/delete/clear를 호출하면 이벤트 루프를 블로킹해 코루틴이 작성된 목적의 동시성을 무너뜨려요. 거기서는 a-접두사 메서드를 사용하세요.
Watcher Trigger 안에서 asset_state_store 사용하기
BaseEventTrigger 서브클래스(watcher trigger)는 run() 안에서 asset state store를 직접 읽고 쓸 수 있어요. Triggerer는 run()이 호출되기 전에 self.asset_state_store를 주입하는데, trigger가 보고 있는 asset에 범위가 지정돼요. __init__이나 serialize() 중에는 사용할 수 없고, run() 안에서만 접근해요.
Task 기반 접근(asset이 inlet·outlet 선언으로 식별되는)과 달리, watcher trigger의 접근자는 자동으로 보고 있는 asset에 바인딩되므로 subscripting이 필요 없어요.
import asyncio
from collections.abc import AsyncIterator
from typing import Any
from airflow.triggers.base import BaseEventTrigger, TriggerEvent
class PollEventsTrigger(BaseEventTrigger):
def __init__(self, source: str, waiter_delay: int, **kwargs):
super().__init__(**kwargs)
self.source = source
self.waiter_delay = waiter_delay
def serialize(self) -> tuple[str, dict[str, Any]]:
return (
f"{self.__class__.__module__}.{self.__class__.__qualname__}",
{"source": self.source, "waiter_delay": self.waiter_delay},
)
def _poll_for_new_record(self, source: str, last_seen: str) -> str | None:
... # Add logic for polling a certain source
return None
async def run(self) -> AsyncIterator[TriggerEvent]:
while True:
last_seen = await self.asset_state_store.aget("last_seen_id", default=0)
new_id = self._poll_for_new_record(
source=self.source,
last_seen=last_seen,
)
if new_id is not None:
await self.asset_state_store.aset("last_seen_id", new_id)
yield TriggerEvent({"status": "success", "record_id": new_id})
return
await asyncio.sleep(self.waiter_delay)
run()은 코루틴이고 triggerer의 모든 trigger는 하나의 이벤트 루프를 공유하므로, 거기서는 비동기 접근자를 사용하세요. 동기 메서드도 작동하지만 API 서버로의 왕복 전체 동안 루프를 붙잡아요.
대응하는 AssetWatcher가 trigger를 asset에 연결해요:
from airflow.sdk import Asset, AssetWatcher
from my_dag.triggers import PollEventsTrigger
my_asset = Asset(
name="orders_api",
watchers=[
AssetWatcher(
name="orders_api_watcher",
trigger=PollEventsTrigger(source="orders", waiter_delay=30),
)
],
)
...
self.asset_state_store는 위 Task 섹션에서 설명한 per-asset 접근자와 동일하게 동작해요: get, set, delete, clear와 그 비동기 대응물 aget, aset, adelete, aclear 모두 사용 가능해요. trigger가 쓴 값은 my_asset을 inlet·outlet으로 선언하는 모든 task가 볼 수 있고 그 반대도 마찬가지예요.
Note
self.asset_state_store는BaseEventTrigger서브클래스 안에서만 사용할 수 있어요. 일반BaseTrigger서브클래스(Task deferral에 사용)는 asset state store에 접근할 수 없어요.
몇 가지 예제 사용 사례
Watermark 패턴
asset state store의 표준 사용 사례는 각 run마다 watermark를 앞으로 진행시키는 증분 로드 task예요. Watermark는 asset 자체에 저장되므로 그 asset을 읽거나 쓰는 모든 task가 접근할 수 있어요. 이 사용 사례는 BaseEventTrigger로 asset "watcher" 같은 것을 만들 때 특히 적합해요.
from airflow.sdk import Asset, DAG, task
orders = Asset("orders", uri="s3://data/orders/")
with DAG("incremental_orders", schedule="@daily"):
@task(inlets=[orders], outlets=[orders])
def load_new_orders(**context):
asset_state_store = context["asset_state_store"] # single-inlet shorthand
# Read the last watermark, default to epoch if first run.
watermark = asset_state_store.get("watermark", default="1970-01-01T00:00:00Z")
# Fetch only rows created after the watermark.
rows = fetch_orders_since(watermark)
if not rows:
return
upload_to_warehouse(rows)
# Advance the watermark to the latest row seen.
new_watermark = max(r["created_at"] for r in rows)
asset_state_store.set("watermark", new_watermark)
각 run마다 task는 이전 run이 남긴 watermark를 읽고, 새 데이터만 가져온 다음 watermark를 진행시켜요. asset state store는 run을 넘어 영속되므로, 재시도, 수동 재실행, Scheduler 재시작에도 다음 run은 이전 run이 멈춘 지점에서 정확히 시작돼요.
수명과 가비지 컬렉션
Asset store 행은 무기한 영속돼요. 이들은 task state store에 적용되는 [state_store] default_retention_days 시간 기반 만료의 대상이 아니에요.
유일한 자동 정리는 *고아 정리(orphan sweep)*예요: asset이 비활성화되면(활성 asset_active 레코드 없음) 그 저장소 행이 다음 가비지 컬렉션 과정에서 제거돼요. 그 정리가 실행되기 전까지는 DB에 오래된 행이 있을 수 있지만 쓰기는 불가능해요. Execution API 해석기는 활성 asset으로만 필터링해요.
asset state store 항목을 명시적으로 제거하려면 Task 안에서 clear()를 호출하거나 REST API를 사용해요.
Note
[state_store] clear_on_success는 asset state store를 지우지 않아요. Asset store는 설계상 여러 run에 걸치므로, 자동 task 레벨 정리는 다음 run이 의존하는 정보를 파괴할 거예요. 더 이상 필요하지 않을 때 asset state store를 항상 명시적으로 지우세요.