Task·Asset 상태 저장소 설정
Task·Asset 상태 저장소 설정 (Task and Asset State Store Configuration)
이 페이지는 버전 3.3에 추가된 Task·Asset 상태 저장소(state store) 설정을 다뤄요. 기본적으로 둘 다 Airflow 메타데이터 데이터베이스에 저장되며, [state_store] 섹션의 옵션(backend, default_retention_days, clear_on_success, state_cleanup_batch_size)과 가비지 컬렉션 의미, 커스텀 백엔드 제공 방법을 설명해요.
출처: 문서
본문
버전 3.3에 추가됨.
Task와 asset 상태 저장소는 task state store와 asset state store의 영속 계층이에요. 기본적으로 둘 다 Airflow 메타데이터 데이터베이스에 저장돼요. 이 페이지에서는 사용 가능한 설정 옵션, 가비지 컬렉션 의미, 커스텀 백엔드를 제공하는 방법을 설명해요. 설정([state_store]), CLI(airflow state-store), 백엔드 기반 클래스(BaseStoreBackend) 모두 이 기능에 state_store 이름을 사용해요.
설정 참조
모든 옵션은 airflow.cfg의 [state_store] 섹션 아래 있어요.
Note
설정 섹션은
[state_store]이며,[task_state_store]가 아닙니다.
backend
BaseStoreBackend를 구현하는 클래스의 전체 dotted 경로. 기본값은 내장 metastore 백엔드예요.
[state_store]
backend = mypackage.state.CustomStateBackend
default_retention_days
task state store 행이 만료된 후 지나는 일수예요. 명시적 보존 기간 없이 키가 쓰여지면, expires_at은 worker에서 now + default_retention_days로 계산돼요. 이 설정을 바꿔도 이미 쓰여진 행에는 영향을 주지 않아요.
0으로 설정하면 시간 기반 정리를 완전히 비활성화해요.- 기본값:
30. - 이 설정은 asset state store 행에는 적용되지 않아요.
[state_store]
default_retention_days = 30
clear_on_success
True이면 task instance가 success 상태로 이동할 때 해당 task instance의 모든 task state store 키가 자동으로 삭제돼요. 기본값은 False로, 관찰 가능성을 위해 success 후에도 task state store 항목을 보존해요 (예: 제출된 job ID나 마지막 row count를 run이 완료된 후에도 UI나 REST API에서 읽을 수 있게요).
Important
clear_on_success는 task state store만 지워요. asset state store에는 영향이 없어요. Asset store는 task instance가 아니라 asset에 범위가 지정되며 명시적으로 지워야 해요.
[state_store]
clear_on_success = False
state_cleanup_batch_size
가비지 컬렉션 정리 중 배치당 삭제되는 행 수예요. 0(기본값)으로 설정하면 일치하는 모든 행을 단일 문으로 삭제해요. 큰 task_state_store 테이블이 있는 배포에서는 잠금 경합을 줄이기 위해 조정해요.
[state_store]
state_cleanup_batch_size = 10000
Worker측 백엔드 ([workers] state_backend)
[workers] 아래의 별도·선택적 설정 키로, task state store와 asset state store 값을 API 서버에 도달하기 전에 worker측 백엔드를 통해 라우팅할 수 있어요.
[workers]
state_backend = mypackage.state.S3StateBackend
이 설정이 있으면 TaskStateStoreAccessor.set()은 값을 Execution API로 보내기 전에 worker측 백엔드에서 serialize_task_state_store_to_ref()를 호출하고(반환된 값은 실제 저장소에 대한 참조), get()은 Execution API에서 저장된 참조를 받은 후 deserialize_task_state_store_from_ref()를 호출해요. 아래 커스텀 worker측 백엔드를 참고해요.
가비지 컬렉션 의미
정리(cleanup) 작업은 "가비지 컬렉션"이라고도 하며 Airflow CLI로 트리거돼요. 정리 작업을 트리거하는 명령은 airflow state-store clean이에요. 이 프로세스는 다음 규칙에 따라 저장소 행을 제거해요:
시간 기반 만료(task state store 전용)
: expires_at < now()인 행이 삭제돼요. expires_at은 서버가 아니라 worker가 쓰기 시점에 계산해요.
default_retention_days 폴백(task state store 전용)
: 명시적 보존 기간 없이 쓰여진 키는 쓰기 시점에 계산된 now + default_retention_days의 expires_at을 얻어요. 가비지 컬렉션은 expires_at < now()인 행을 삭제해요.
NEVER_EXPIRE 키
: retention=NEVER_EXPIRE로 설정된 키는 expires_at = NULL과, 가비지 컬렉션이 그것을 무조건 건너뛰도록 알려주는 플래그로 저장돼요. default_retention_days와 무관하게 시간 기반 정리로 절대 삭제되지 않아요.
on_delete=CASCADE (asset state store)
: Asset이 삭제되면 그 asset에 대한 모든 asset state store 행이 삭제돼요.
Important
가비지 컬렉션은
MetastoreBackend에서만 작동해요. 커스텀 백엔드는 명시적으로 건너뛰어져요.
커스텀 백엔드
커스텀 백엔드는 BaseStoreBackend를 서브클래싱하고 그 추상 메서드를 구현해야 해요: 동기 호출자용 get, set, delete, clear와 aget, aset, adelete, aclear 비동기 대응물. 전체 API는 BaseStoreBackend를 참고해요.
각 메서드는 TaskScope 또는 AssetScope인 scope 인자를 받아요. isinstance를 사용해 디스패치해요:
from airflow.sdk.state import BaseStoreBackend, TaskScope, AssetScope
class MyBackend(BaseStoreBackend):
def get(self, scope, key, *, session=None):
if isinstance(scope, TaskScope):
return self._task_store.get(scope, key)
elif isinstance(scope, AssetScope):
return self._asset_store.get(scope, key)
AssetScope에는 세 개의 선택적 필드가 있어요: asset_id(정수, 서버측 전용), name, uri. 적어도 하나는 설정해야 해요. 서버측 작업(REST API 호출)은 asset_id를 제공해요. worker측 작업은 name이나 uri를 제공해요(worker는 정수 asset_id에 접근할 수 없어요).
클래스를 [state_store] backend를 통해 구성해요:
[state_store]
backend = mypackage.state.MyBackend
커스텀 worker측 백엔드
Worker측 백엔드는 두 쌍의 직렬화 훅으로 BaseStoreBackend를 확장해요. 이들은 [workers] state_backend로 별도 구성되며 API 서버가 아니라 worker 프로세스에서 실행돼요. 이를 통해 큰 페이로드나 자격 증명 데이터를 worker 인프라를 사용해 직접 저장하고, 데이터베이스에는 압축된 참조 문자열만 유지할 수 있어요.
BaseStoreBackend의 네 개 직렬화 훅을 오버라이드해요:
serialize_task_state_store_to_ref: 값이 Execution API로 보내지기 전에TaskStateStoreAccessor.set()이 호출해요. 원시 값 대신 데이터베이스에 저장할 압축 참조 문자열(예: S3 키)을 반환해요.deserialize_task_state_store_from_ref: 백엔드에서 참조를 검색한 후TaskStateStoreAccessor.get()이 호출해요. 실제 값을 반환해요.serialize_asset_state_store_to_ref: task 변형과 같지만 asset state store용이에요.scope로 asset scope(name및/또는uri를 가진AssetScope)를 받아요.deserialize_asset_state_store_from_ref: 저장된 참조를 실제 값으로 해석하기 위해AssetStateStoreAccessor.get()이 호출해요.
Important
참조는 결정적(deterministic)이어야 합니다. 같은 입력(
scope+key)이 주어지면 직렬화 메서드는 항상 같은 참조 문자열을 반환해야 해요. 참조 경로에 타임스탬프, 랜덤 UUID, 그 외 비결정적 구성 요소를 넣지 마세요.키가 삭제되거나 지워지면 Airflow는 데이터베이스 참조를 먼저 지운 다음 백엔드의
delete()나clear()메서드를 호출해요. DB 행이 사라진 후 백엔드 정리가 실패하면 외부 객체는 고아(orphaned)가 돼요. 참조가 결정적이므로, 같은 키에 대한 이후set()이 그 고아 객체를 덮어써서 실패를 복구 가능하게 해요. 비결정적 참조는 외부 객체를 영구적으로 고아로 남겨 찾을 방법이 없게 해요.
예제 스켈레톤:
from airflow.sdk.state import BaseStoreBackend, TaskScope, AssetScope
if TYPE_CHECKING:
from pydantic import JsonValue
class S3StateBackend(BaseStoreBackend):
def _task_ref(self, scope: TaskScope, key: str) -> str:
return f"airflow/task-store/{scope.dag_id}/{scope.run_id}/{scope.task_id}/{scope.map_index}/{key}"
def _asset_ref(self, scope: AssetScope, key: str) -> str:
import hashlib
asset_identifier = scope.name or scope.uri or ""
safe = hashlib.sha256(asset_identifier.encode()).hexdigest()[:16]
return f"airflow/asset-store/{safe}/{key}"
def serialize_task_state_store_to_ref(self, *, value: JsonValue, key: str, scope: TaskScope) -> str:
s3_key = self._task_ref(scope, key)
s3_client.put_object(Bucket=BUCKET, Key=s3_key, Body=json.dumps(value).encode())
return s3_key
def deserialize_task_state_store_from_ref(self, stored: str) -> JsonValue:
s3_object = s3_client.get_object(Bucket=BUCKET, Key=stored)
return json.loads(s3_object["Body"].read().decode())
def serialize_asset_state_store_to_ref(self, *, value: JsonValue, key: str, scope: AssetScope) -> str:
s3_key = self._asset_ref(scope, key)
s3_client.put_object(Bucket=BUCKET, Key=s3_key, Body=json.dumps(value).encode())
return s3_key
def deserialize_asset_state_store_from_ref(self, stored: str) -> JsonValue:
s3_object = s3_client.get_object(Bucket=BUCKET, Key=stored)
return json.loads(s3_object["Body"].read().decode())
# Implement the remaining abstract methods as pass-throughs or delegating to the
# default MetastoreBackend for the DB side
...