Asset 이벤트 발행·소비와 AssetAlias, 크로스팀 접근제어
Apache Airflow Official Documentation (in-repo snapshot) — Apache Software Foundation authoring-and-scheduling/assets.rst - Creating a task to emit asset events, Fetching information from previously emitted asset events, Output to multiple assets, Dynamic data events emitting and asset creation through AssetAlias, Cross-team asset event filtering with producer_teams
이 모듈을 다 읽으면
- outlets/inlets로 Asset 이벤트를 발행·소비하는 방법과 @asset/@asset.multi 데코레이터가 만드는 리소스를 설명할 수 있다
- inlet_events 접근자의 체이닝 필터를 활용해 원하는 이전 이벤트를 조회할 수 있다
- AssetAlias가 필요한 상황과 동일 (asset, extra) 조합이 중복 이벤트를 만들지 않는다는 규칙을 설명할 수 있다
- Multi-Team 모드에서 producer_teams/consumer_teams/allow_global이 이벤트 전달에 미치는 영향을 판단할 수 있다
태스크는 outlets로 Asset 이벤트를 발행하고 inlets/inlet_events로 이전 이벤트를 읽는다. @asset·@asset.multi 데코레이터는 이 boilerplate를 줄여주는 shorthand다. URI가 실행 시점에야 정해지는 경우에는 AssetAlias로 동적으로 Asset을 만들 수 있고, Multi-Team 모드에서는 producer_teams·consumer_teams·allow_global로 어느 팀의 이벤트가 어느 팀의 Dag를 트리거할 수 있는지 세밀하게 통제한다.
outlets로 Asset 이벤트 발행하기, @asset 데코레이터
태스크는 ``outlets`` 인자에 Asset을 지정해 이벤트를 발행한다.
from airflow.sdk import DAG, Asset
from airflow.providers.standard.operators.python import PythonOperator
example_asset = Asset(name="example_asset", uri="s3://asset-bucket/example.csv")
def _write_example_asset():
...
with DAG(dag_id="example_asset", schedule="@daily"):
PythonOperator(task_id="example_asset", outlets=[example_asset], python_callable=_write_example_asset)
'하나의 태스크가 하나의 Asset에 이벤트를 발행하는 Dag'라는 가장 흔한 패턴을 위해 Airflow는 ``@asset`` 데코레이터라는 shorthand를 제공한다. 위 코드와 완전히 동일한 결과를 만드는 코드는 다음과 같다.
from airflow.sdk import asset
@asset(uri="s3://asset-bucket/example.csv", schedule="@daily")
def example_asset():
...
``@asset``을 선언하면 함수 이름을 name으로 하는 ``Asset``, 함수 이름을 dag_id로 하는 ``DAG``, 그리고 함수 이름을 task_id로 하고 그 Asset을 outlet으로 갖는 태스크가 자동으로 만들어진다. ``@asset`` 함수 안에서 ``self``, ``context``, ``outlet_events``라는 파라미터 이름은 예약어로, Airflow가 실행 시점에 각각 해당 Asset 자신·실행 컨텍스트·outlet 이벤트 접근자를 채워 넣으며, 이 이름들은 inlet Asset 참조로 취급되지 않는다.
``@asset``과 ``@task``를 함께 쓰면(예: ``@asset(...)`` 아래 ``@task.bash(retries=3)``) TaskFlow 방식으로 태스크의 초기 인자를 지정하거나 BashOperator 같은 다른 오퍼레이터를 활용하는 식으로 커스터마이즈할 수 있다.
핵심 포인트
- 태스크는 outlets=[asset]으로 Asset 이벤트를 발행하며, @asset 데코레이터는 '태스크 1개가 Asset 1개를 발행하는 Dag' 패턴의 boilerplate를 줄여준다
- @asset은 함수 이름으로 Asset·DAG·태스크를 자동 생성하며, self/context/outlet_events는 예약된 파라미터 이름이다
Asset 이벤트에 extra 정보 첨부하기
Asset 자체의 extra(정적 메타데이터)와 달리, Asset '이벤트'의 extra는 그 갱신을 촉발한 데이터 변화를 주석하는 데 쓴다 — 예를 들어 이번 갱신으로 몇 개의 행이 바뀌었는지, 어떤 날짜 범위를 커버하는지 같은 정보다. 가장 쉬운 방법은 태스크에서 ``Metadata`` 객체를 yield하는 것이다.
from airflow.sdk import Metadata, asset
@asset(uri="s3://asset/example.csv", schedule=None)
def example_s3(self):
df = ...
yield Metadata(self, {"row_count": len(df)})
또는 태스크 실행 컨텍스트의 ``outlet_events`` 접근자에 직접 대입해도 된다: ``context["outlet_events"][self].extra = {"row_count": len(df)}``. 클래식 오퍼레이터에서는 ``execute``를 오버라이드하는 서브클래싱이 권장되는 방법이며, ``pre_execute``/``post_execute`` 훅에서도 가능하지만 이 훅들은 태스크가 재시도될 때 다시 실행되지 않으므로 extra 정보가 실제 데이터와 어긋날 수 있다는 점을 기억해야 한다. extra 정보는 DB에 저장되므로 JSON 직렬화 가능한 값(list·dict 중첩 포함)만 담을 수 있다.
핵심 포인트
- Asset 자체의 extra는 정적 메타데이터, Asset 이벤트의 extra는 그 갱신을 촉발한 데이터 변화(행 수, 날짜 범위 등)를 담는 용도로 성격이 다르다
- Metadata를 yield하거나 outlet_events[self].extra에 직접 대입해 이벤트 extra를 채울 수 있으며, pre_execute/post_execute 훅은 재시도 시 재실행되지 않아 extra가 실제와 어긋날 위험이 있다
inlets로 이전 Asset 이벤트 조회하기
outlets로 발행된 이벤트는 같은 Asset을 ``inlets``에 선언한 태스크가 ``inlet_events`` 접근자로 읽을 수 있다. 각 값은 timestamp 순(오래된 것부터)으로 정렬된 시퀀스이며 ``[-1]``, ``[-2:]``처럼 파이썬 리스트 인터페이스 대부분을 지원한다. 이 접근자는 지연 평가(lazy)되어 실제로 항목에 접근할 때만 DB를 조회한다.
@asset(schedule=None)
def post_process_s3_file(context, write_to_s3):
events = context["inlet_events"][write_to_s3]
last_row_count = events[-1].extra["row_count"]
체이닝 필터도 지원한다. ``.partition_key(value)``는 파티션 키 정확 일치, ``.partition_key_regexp_pattern(pattern)``은 정규식 필터링인데 후자는 ``[api] regexp_query_timeout``을 양수로 설정해야만 활성화되는 opt-in 기능이다(쿼리 실행 시간도 이 값으로 제한된다). ``.extra(key, value)``는 extra 키-값 조건으로 필터링하며 여러 번 체이닝하면 AND로 결합된다. 그 외 ``.after(timestamp)``, ``.before(timestamp)``, ``.ascending()``, ``.limit(n)``도 체이닝할 수 있다. REST API에서도 ``extra`` 쿼리 파라미터를 반복 지정해 같은 필터링을 할 수 있다 (``key=value`` 형식, 여러 조건은 AND로 결합).
핵심 포인트
- inlet_events는 이전 Asset 이벤트를 timestamp 오름차순으로 담은 지연 평가 시퀀스이며, [-1] 같은 리스트 인터페이스를 지원한다
- partition_key_regexp_pattern은 [api] regexp_query_timeout을 양수로 설정해야 활성화되는 opt-in 기능이고, .extra(key, value) 체이닝은 AND로 결합된다
@asset.multi로 여러 Asset에 출력하기
하나의 태스크가 여러 Asset에 이벤트를 발행해야 하는 경우(예: 하나의 데이터 소스를 여러 갈래로 나눠야 할 때)도 있다. 일반적으로는 권장되지 않지만 필요할 때는 ``outlets``를 복수로 지정하면 된다.
@task(inlets=[input_asset], outlets=[out_asset_1, out_asset_2])
def process_input():
...
이를 위한 shorthand가 ``@asset.multi``이며, ``outlets``를 데코레이터 인자로 지정하고 함수 인자로는 inlet Asset을 받는다.
@asset.multi(schedule=None, outlets=[out_asset_1, out_asset_2])
def process_input(input_asset):
...
핵심 포인트
- 한 태스크가 여러 Asset에 이벤트를 발행하는 것은 일반적으로 권장되지 않지만, outlets를 복수로 주거나 @asset.multi shorthand로 구현할 수 있다
AssetAlias로 동적 Asset 생성하기
URI 같은 Asset의 고정 속성이 실행 전에는 알려지지 않는 경우 ``AssetAlias``를 쓴다. 태스크는 alias를 안정적인 이름으로 ``outlets``에 선언해두고, 실행 시점에 ``outlet_events``나 yield된 ``Metadata``를 통해 하나 이상의 실제 ``Asset``을 alias에 연결한다.
from airflow.sdk import AssetAlias
@task(outlets=[AssetAlias("my-task-outputs")])
def my_task_with_outlet_events(*, outlet_events):
outlet_events[AssetAlias("my-task-outputs")].add(Asset("s3://bucket/my-task"), extra={"k": "v"})
동일한 (asset, extra) 조합은 같은 alias에 여러 번 add하거나 여러 alias에 add해도 이벤트가 한 번만 발행된다. 반대로 extra 값이 다르면 추가 이벤트가 발행된다 — 예를 들어 같은 Asset을 세 개의 alias에 각각 add하되 첫 두 번은 ``{"k": "v"}``로 동일하고 세 번째만 ``{"k2": "v2"}``라면, 이벤트는 총 2개만 발행된다.
다운스트림 Dag는 alias 이름을 그대로 참조해 ``schedule``이나 ``inlets``에 쓸 수 있고, alias가 실제로 어떤 Asset으로 해석되는지는 신경 쓰지 않아도 된다. alias를 통해 트리거된 이벤트를 읽을 때도 ``inlet_events[AssetAlias(...)]`` 형태로 동일하게 접근한다.
핵심 포인트
- AssetAlias는 URI 같은 Asset의 고정 속성이 실행 전에 정해지지 않을 때, 안정적인 이름을 outlets에 선언해두고 실행 시점에 실제 Asset을 연결하는 방식이다
- 동일한 (asset, extra) 조합은 같은 alias나 여러 alias에 반복 add해도 이벤트가 한 번만 발행되고, extra가 다르면 추가 이벤트가 발행된다
크로스팀 Asset 이벤트 필터링 (producer_teams)
Multi-Team 모드가 켜져 있으면 Asset 이벤트는 팀 소속에 따라 필터링된다. 기본적으로 소비 Dag는 같은 팀 소속이거나 팀이 없는(글로벌) Dag가 발행한 이벤트만 받는다 — 의도치 않은 팀 간 트리거를 막기 위해서다. 교차 팀 접근을 허용하려면 Asset 정의에 ``AssetAccessControl``을 지정한다.
from airflow.sdk import Asset, AssetAccessControl
shared_data = Asset(
name="my_data",
uri="s3://bucket/shared/data.csv",
access_control=AssetAccessControl(producer_teams=["team_analytics", "team_ml"]),
)
``AssetAccessControl``은 producer_teams(이 Asset의 소비자가 자기 팀 이벤트 외에 추가로 받아들일 발행 팀 목록, 기본 빈 리스트), consumer_teams(이 Asset의 이벤트를 소비할 수 있는 팀 목록, 기본 None), allow_global(팀 없는 글로벌 Dag가 교차 팀 이벤트 전달에 참여할지 여부, 기본 True) 세 파라미터를 받는다. allow_global=False로 설정하면 글로벌(팀 없는) producer Dag가 이 Asset의 소비자를 트리거하지 못하게 막을 수 있다(단, allow_global은 Dag producer에만 영향을 주고, 팀 없는 API 사용자는 이 설정과 무관하게 항상 팀 없는 소비자만 트리거할 수 있다).
access_control을 지정하지 않으면 기본 규칙이 적용된다: 양쪽이 같은 팀이면 항상 전달, producer가 팀이 있고 consumer가 다른 팀이면 producer_teams에 없는 한 차단, producer가 글로벌이면 allow_global=True인 소비자 전원에게 전달(기본값), consumer가 글로벌이면 어떤 producer의 이벤트든 수용, 양쪽 다 팀이 없으면 전달. Multi-Team 모드가 꺼져 있으면 access_control은 무시되고 모든 이벤트가 모든 소비 Dag에 전달되어 이전 동작과 호환된다.
핵심 포인트
- Multi-Team 모드에서 소비 Dag는 기본적으로 같은 팀 또는 글로벌 Dag의 이벤트만 받으며, AssetAccessControl(producer_teams, consumer_teams, allow_global)로 교차 팀 허용 범위를 조정한다
- allow_global=False는 글로벌 producer Dag의 트리거를 차단하지만, 팀 없는 API 사용자는 이 설정과 무관하게 항상 팀 없는 소비자만 트리거할 수 있다
- Multi-Team 모드가 비활성화되어 있으면 access_control 설정 자체가 무시되고 모든 이벤트가 모든 소비 Dag에 전달된다