Lineage

Lineage

이 페이지는 Airflow의 데이터 계보(lineage) 추적 기능을 다뤄요. Task 사이뿐 아니라 Task 안에서 사용하는 훅(hook)으로부터 데이터가 어떻게 흘러가는지 이해할 수 있게 도와주는 강력한 기능이에요. 다만 아직 실험 단계라 변경될 수 있다는 점을 알아두세요. HookLineageCollector가 계보 정보를 수집하는 중심 역할을 하고, HookLineageReader로 그 정보를 읽어올 수 있어요.

출처: 문서

본문

Note

Lineage 지원은 매우 실험적이며 변경될 수 있어요.

Airflow는 Task 사이뿐 아니라 그 Task 안에서 사용되는 훅으로부터도 데이터 계보(lineage)를 추적하는 강력한 기능을 제공해요. 이 기능은 데이터가 Airflow 파이프라인을 통해 어떻게 흘러가는지를 이해하는 데 도움을 줘요.

HookLineageCollector의 전역 인스턴스가 계보 정보를 수집하는 중심 허브 역할을 해요. 훅은 자신이 상호작용하는 Asset에 대한 상세 정보를 이 콜렉터로 보낼 수 있어요. 콜렉터는 이 데이터를 사용해 Asset을 설명하는 표준 형식인 AIP-60 호환 Asset을 구성해요. 훅은 아래 예제처럼 Asset과 무관한 임의의 데이터도 이 콜렉터로 보낼 수 있어요.

from airflow.sdk.lineage import get_hook_lineage_collector

class CustomHook(BaseHook):
    def run(self):
        # run actual code
        collector = get_hook_lineage_collector()
        collector.add_input_asset(self, asset_kwargs={"scheme": "file", "path": "/tmp/in"})
        collector.add_output_asset(self, asset_kwargs={"scheme": "file", "path": "/tmp/out"})
        collector.add_extra(self, key="external_system_job_id", value="some_id_123")

HookLineageCollector가 수집한 계보 데이터는 Airflow 플러그인에 등록된 HookLineageReader 인스턴스를 통해 접근할 수 있어요.

from airflow.sdk.lineage import HookLineageReader
from airflow.plugins_manager import AirflowPlugin

class CustomHookLineageReader(HookLineageReader):
    def get_inputs(self):
        return self.lineage_collector.collected_assets.inputs

class HookLineageCollectionPlugin(AirflowPlugin):
    name = "HookLineageCollectionPlugin"
    hook_lineage_readers = [CustomHookLineageReader]

Airflow 안에 HookLineageReader가 등록되어 있지 않으면 기본 NoOpCollector가 대신 사용돼요. 이 콜렉터는 AIP-60 호환 Asset을 만들지도, 계보 정보를 수집하지도 않아요.

더 알아보기 (Learn more)