Airflow 3.0+ 공개 인터페이스

Airflow 3.0+ 공개 인터페이스 (Public Interface for Airflow 3.0+)

Apache Airflow의 Public Interface는 변경이 시맨틱 버저닝에 의해 관리되는 인터페이스와 동작의 모음이에요. 이 문서는 Airflow 3의 공개 인터페이스와 그 사용 방법, 그리고 공개 인터페이스가 아닌 것에 대해 설명해 드려요.

출처: 문서

본문

Warning

이 문서는 Airflow 3.0+의 Public Interface를 다룹니다.

Airflow 2.x를 사용 중이라면 레거시 인터페이스에 대해 Airflow 2.11 Public Interface 문서를 참고하세요.

Apache Airflow의 Public Interface는 변경이 시맨틱 버저닝에 의해 관리되는 Apache Airflow의 인터페이스와 동작의 모음이에요. 사용자는 Dags를 만들고 관리하고, tasks와 의존성을 관리하고, 새 executors, plugins, operators, providers를 작성해 Airflow 기능을 확장함으로써 Airflow의 public interface와 상호작용해요. Public Interface는 커스텀 도구와 다른 시스템과의 통합을 구축하고, Airflow 워크플로우의 특정 측면을 자동화하는 데 유용할 수 있어요.

Dag 작성자와 task 실행을 위한 기본 공개 인터페이스는 task SDK를 사용하는 거예요. Airflow task SDK는 Dag 작성자와 task 실행을 위한 기본 공개 인터페이스예요 airflow.sdk 네임스페이스. task 코드에서 메타데이터 데이터베이스에 직접 접근하는 것은 더 이상 허용되지 않아요. 대신 Stable REST API, Python Client, 또는 Task Context 메서드를 사용하세요.

포괄적인 Task SDK 문서는 Task SDK 참조를 참고하세요.

Airflow Public Interfaces 사용하기

Note

Airflow 3.0부터 사용자는 AIP-72에 정의된 대로 airflow.sdk 네임스페이스를 공식 Public Interface로 사용해야 해요.

내부 모듈이나 메타데이터 데이터베이스와의 직접 상호작용은 불가능해요. 안정적이고 프로덕션에 안전한 통합을 위해 다음을 권장해요:

  • 공식 REST API
  • Python Client SDK (airflow-client-python)
  • Task SDK (airflow.sdk)

관련 문서:

다음은 Airflow public interface의 몇 가지 예시예요:

  • 자신만의 operators나 hooks를 작성할 때. 이는 사용 사례에 맞는 hook이나 operator가 없을 때, 또는 있지만 동작을 커스터마이즈해야 할 때 흔히 합니다.
  • Dag 빌딩 블록을 넘어 Airflow의 기능을 확장하는 새 Plugins를 작성할 때. Secrets, Timetables, Triggers, Listeners가 모두 그러한 기능의 예시예요. 이는 보통 Airflow 인스턴스를 관리하는 사용자가 해요.
  • 커스텀 Operators, Hooks, Plugins를 번들링하고 providers로 함께 릴리스할 때 — 보통 Airflow가 통합하는 외부 서비스나 애플리케이션을 위한 재사용 가능한 기능 세트를 제공하려는 사람이 해요.
  • task를 작성하기 위해 taskflow API를 사용하기
  • Airflow 객체의 일관된 동작에 의존하기

"public interface"의 한 측면은 Airflow Python 클래스와 함수를 확장하거나 사용하는 것이에요. 아래 언급된 클래스와 함수는 Airflow의 MAJOR 버전 내에서 역호환 가능한 시그니처와 동작을 유지한다고 믿을 수 있어요. 반면 _로 시작하는 클래스와 메서드(보호된 Python 메서드로도 알려짐)와 __로 시작하는 것(비공개 Python 메서드로도 알려짐)은 Public Airflow Interface의 일부가 아니며 언제든 변경될 수 있어요.

Stable REST API(OpenAPI 사양 기반)를 통해 Airflow의 Public Interface를 사용할 수도 있어요. 특정 요구 사항에 따라 Airflow Command Line Interface (CLI)를 사용할 수도 있지만, 그 동작은 세부 사항(출력 형식과 사용 가능한 플래그 등)에서 변경될 수 있으므로, 이를 프로그래밍 방식으로 의존하려면 Stable REST API를 권장해요.

Dag 작성자를 위한 Public Interface 사용하기

Dag 작성자의 기본 인터페이스는 airflow.sdk 네임스페이스예요. 포괄적인 문서는 Task SDK 참조를 참고하세요. 이것은 내부 구현 변경의 대상이 아닌, Dags와 tasks를 만들기 위한 안정적이고 잘 정의된 인터페이스를 제공해요. 이 변경의 목표는 Dag 작성을 Airflow 내부(Scheduler, API Server 등)에서 분리해, Airflow 버전을 넘나드는 Dags 작성과 유지를 위한 버전에 구애받지 않는 안정적 인터페이스를 제공하는 것이에요.

airflow.sdk의 주요 Imports:

Classes:

Decorators and Functions:

참고 (See also)

TaskGroup, DAG, task_group의 API 참조

Airflow 2.x에서의 마이그레이션:

Airflow 2.x에서 3.x로의 자세한 마이그레이션 지침(import 변경과 다른 breaking changes 포함)은 마이그레이션 가이드를 참고하세요.

사용 가능한 클래스, decorators, 함수의 전체 목록은 airflow.sdk.__all__을 확인하세요.

모든 Dags는 내부 Airflow 모듈을 직접 참조하는 대신 airflow.sdk를 사용하도록 imports를 업데이트해야 해요. 레거시 import 경로(예: airflow.models.dag.DAG, airflow.decorator.task)는 폐지되며 향후 Airflow 버전에서 제거될 거예요.

Dags

Dag은 반복되는 워크플로우를 나타내는 Airflow의 핵심 엔티티예요. Dag 파일에서 DAG 클래스를 인스턴스화해 Dag을 만들 수 있어요. Dags는 Param 클래스로 파라미터를 지정할 수도 있어요.

Dags를 만드는 권장 방법은 airflow.sdk 네임스페이스의 dag() decorator를 사용하는 거예요.

Airflow에는 Dags 작성 방법을 배울 수 있는 예시 Dags 세트가 있어요.

Dags에 대한 자세한 내용은 Dags에서 읽을 수 있어요.

Dags에서 사용되는 모듈의 참조는 여기에 있어요:

Note

airflow.sdk 네임스페이스는 Dag 작성자를 위한 기본 인터페이스를 제공해요. 자세한 API 문서는 Task SDK 참조를 참고하세요.

Note

DagBag 클래스는 Airflow가 파일과 폴더에서 Dags를 로드하는 데 내부적으로 사용돼요. Dag 작성자는 대신 airflow.sdk 네임스페이스의 DAG 클래스를 사용해야 해요.

Note

DagRun 클래스는 Airflow가 Dag run 관리를 위해 내부적으로 사용해요. Dag 작성자는 get_current_context()를 통한 Task Context를 통해 Dag run 정보에 접근하거나 DagRunProtocol 인터페이스를 사용해야 해요.

Operators

기본 클래스 BaseOperatorBaseSensorOperator는 공개되며 새 operators를 만들도록 확장할 수 있어요.

새 operators의 기본 클래스는 airflow.sdk 네임스페이스의 BaseOperator예요.

Apache Airflow에서 게시된 BaseOperator의 하위 클래스는 구조가 아니라 동작에서 공개돼요. 즉 Operator의 파라미터와 동작은 semver에 의해 관리되지만 메서드는 언제든 변경될 수 있어요.

Task Instances

Task instances는 (Dag Run에서) Dag의 단일 task의 개별 실행이에요. Task instances는 get_current_context()를 통해 Task Context로 접근돼요. 직접 데이터베이스 접근은 불가능해요.

Note

Task Context는 airflow.sdk 네임스페이스의 일부예요. 자세한 API 문서는 Task SDK 참조를 참고하세요.

Task Instance Keys

Task instance keys는 (Dag Run에서) Dag의 task 인스턴스의 고유 식별자예요. 키는 dag_id, task_id, run_id, try_number, map_index로 구성된 튜플이에요.

task 코드에서 TaskInstance 모델을 통한 task instance keys에 대한 직접 접근은 더 이상 허용되지 않아요. 대신 task instance 정보에 접근하기 위해 get_current_context()를 통한 Task Context를 사용하세요.

Task Context를 통해 task instance 정보에 접근하는 예시:

from airflow.sdk import get_current_context


def my_task():
    context = get_current_context()
    ti = context["ti"]

    dag_id = ti.dag_id
    task_id = ti.task_id
    run_id = ti.run_id
    try_number = ti.try_number
    map_index = ti.map_index

    print(f"Task: {dag_id}.{task_id}, Run: {run_id}, Try: {try_number}, Map Index: {map_index}")

Note

TaskInstanceKey 클래스는 Airflow가 task 인스턴스를 식별하기 위해 내부적으로 사용해요. Dag 작성자는 대신 get_current_context()를 통한 Task Context로 task instance 정보에 접근해야 해요.

Hooks

Hooks는 가능하면 공통 인터페이스를 구현하고 operators의 빌딩 블록 역할을 하는 외부 플랫폼과 데이터베이스에 대한 인터페이스예요. 모든 hooks는 BaseHook에서 파생돼요.

Airflow에는 공개로 간주되는 Hooks 세트가 있어요. 그것들을 확장해 기능을 자유롭게 넓힐 수 있어요:

공개 Airflow 유틸리티 (Public Airflow utilities)

Hooks와 Operators를 작성하거나 확장할 때 Dag 작성자와 개발자는 다음 클래스를 사용할 수 있어요:

  • 외부 서비스 자격 증명과 구성을 제공하는 Connection.
  • Airflow 구성 변수를 제공하는 Variable.
  • task 간 통신 데이터에 접근하는 데 사용되는 XCom.

Connection 및 Variable 작업은 get_current_context()와 task instance의 메서드를 통한 Task Context를 통해, 또는 airflow.sdk 네임스페이스를 통해 수행해야 해요. task 코드에서 ConnectionVariable 모델에 대한 직접 데이터베이스 접근은 더 이상 허용되지 않아요.

Task Context를 통해 Connections와 Variables를 접근하는 예시:

from airflow.sdk import get_current_context


def my_task():
    context = get_current_context()

    conn = context["conn"]
    my_connection = conn.get("my_connection_id")

    var = context["var"]
    my_variable = var.value.get("my_variable_name")

airflow.sdk 네임스페이스를 직접 사용하는 예시:

from airflow.sdk import Connection, Variable

conn = Connection.get("my_connection_id")
var = Variable.get("my_variable_name")

공개 Airflow 유틸리티에 대한 자세한 내용은 Managing Connections, Variables, XComs에서 읽을 수 있어요.

유틸리티에 사용되는 클래스의 참조는 여기 있어요:

Note

Connection, Variable, XCom 클래스는 이제 airflow.sdk 네임스페이스의 일부예요. 자세한 API 문서는 Task SDK 참조를 참고하세요.

공개 예외 (Public Exceptions)

커스텀 Operators와 Hooks를 작성할 때 Airflow가 노출하는 공개 Exceptions를 처리하고 발생시킬 수 있어요:

공개 유틸리티 클래스 (Public Utility classes)

Public Interface를 사용해 Airflow 기능 확장하기

Airflow는 Plugin 메커니즘을 사용해 Airflow 플랫폼 기능을 확장해요. Airflow UI를 확장할 수 있게 하면서도 아래 커스터마이제이션(Triggers, Timetables, Listeners 등)을 노출하는 방식이에요. Providers도 plugin 엔드포인트를 구현하고 Airflow UI와 커스터마이제이션을 커스터마이즈할 수 있어요.

플러그인에 대한 자세한 내용은 Plugins에서 읽을 수 있어요. Apache 웹 UI에서 뷰 커스터마이즈에서 Airflow UI를 확장하는 방법을 읽을 수 있어요. 플러그인이 필요 없는 몇 가지 간단한 UI 커스터마이제이션이 있다는 점에 유의하세요 — UI 커스터마이징에서 이에 대한 자세한 내용을 읽을 수 있어요.

Plugins로 Airflow를 확장할 수 있는 방법은 다음과 같아요:

Triggers

Airflow는 asyncio 호환 Deferrable Operators를 구현하기 위해 Triggers를 사용해요. 모든 Triggers는 BaseTrigger에서 파생돼요.

Airflow에는 공개로 간주되는 Triggers 세트가 있어요. 그것들을 확장해 기능을 자유롭게 넓힐 수 있어요:

Triggers에 대한 자세한 내용은 Deferrable Operators & Triggers에서 읽을 수 있어요.

Timetables

커스텀 timetable 구현은 Airflow 스케줄러에 내장 스케줄 표현으로는 불가능한 방식으로 Dag runs를 스케줄하는 추가 로직을 제공해요. 모든 Timetables는 Timetable에서 파생돼요.

Airflow에는 공개로 간주되는 Timetables 세트가 있어요. 그것들을 확장해 기능을 자유롭게 넓힐 수 있어요:

Timetables에 대한 자세한 내용은 Timetables로 Dag 스케줄링 커스터마이즈에서 읽을 수 있어요.

Listeners

Listeners를 사용하면 Dag/Task 라이프사이클 이벤트에 응답할 수 있어요.

이것은 Dag/Task 라이프사이클 이벤트에 응답하기 위해 구현할 수 있는 hooks를 제공하는 ListenerManager 클래스로 구현돼요.

버전 2.5에 추가: Listener 공개 인터페이스는 버전 2.5에 추가됨.

Listeners에 대한 자세한 내용은 Listeners에서 읽을 수 있어요.

Extra links는 커스텀 Operators와 독립적으로 Airflow에 추가할 수 있는 동적 링크예요. 보통 Operators로 정의할 수 있지만, 플러그인은 링크를 전역 수준에서 재정의할 수 있게 해줘요.

Extra Links에 대한 자세한 내용은 Operator extra link 정의에서 읽을 수 있어요.

Public Interface를 사용해 외부 서비스·애플리케이션과 통합하기

Airflow의 Tasks는 Hooks와 Operators를 통해 외부 서비스를 오케스트레이션할 수 있어요. Airflow의 핵심 기능(예: 인증)도 외부 서비스를 활용하도록 확장할 수 있어요. providers에 대한 자세한 내용은 providers와 그것들이 제공할 수 있는 core extensions을 providers에서 읽을 수 있어요.

Executors

Executors는 task 인스턴스가 실행되는 메커니즘이에요. 모든 executors는 BaseExecutor에서 파생돼요. Airflow에 내장된 executor 구현이 몇 가지 있으며, 각각 고유한 특성과 기능이 있어요.

executor 인터페이스 자체(BaseExecutor 클래스)는 공개지만, 내장 executors(KubernetesExecutor, LocalExecutor 등)는 공개가 아니에요. 즉 예를 들어 KubernetesExecutor를 사용한다면, KubernetesExecutor를 서브클래싱하는 executor를 깨뜨릴 수 있는 변경을 KubernetesExecutor에 minor 또는 patch Airflow 릴리스에서 할 수 있어요. 이것은 Airflow 개발자들이 우리가 제공하는 executors를 계속 개선할 충분한 자유를 허용하기 위해 필요해요. 따라서 내장 executor를 수정하거나 확장하려면 전체 executor 코드를 프로젝트에 통합해 그런 변경이 파생 executor를 깨뜨리지 않도록 해야 해요.

executors와 나만의 것을 작성하는 방법에 대한 자세한 내용은 Executor에서 읽을 수 있어요.

버전 2.6에 추가: executor 인터페이스는 한동안 Airflow에 존재해 왔지만, 2.6 이전에는 코드베이스의 다른 곳에 executor별 코드가 있었어요. 버전 2.6부터 executors는 완전히 분리되어, Airflow core가 특정 executors의 동작을 알 필요가 없어졌어요. Airflow 2.6 이전에도 커스텀 executor 구현에 성공할 수 있었고 몇몇 사람이 그랬지만, 내장 executors를 선호하는 몇 가지 하드코딩된 동작이 있었고, 커스텀 executors는 내장 executors가 가진 완전한 기능을 제공할 수 없었어요.

Secrets Backends

Airflow는 ConnectionVariable을 검색하기 위해 secrets backends에 의존하도록 구성할 수 있어요. 모든 secrets backends는 BaseSecretsBackend에서 파생돼요.

모든 Secrets Backend 구현은 공개예요. 기능을 확장할 수 있어요:

Secret Backends에 대한 자세한 내용은 Secrets Backend에서 읽을 수 있어요. 커뮤니티 providers에서 구현된 사용 가능한 모든 Secrets Backends는 Secret backends에서 찾을 수 있어요.

Auth managers

Auth managers는 Airflow에서 사용자 인증과 사용자 권한 부여를 담당해요. 모든 auth managers는 BaseAuthManager에서 파생돼요.

auth manager 인터페이스 자체(BaseAuthManager 클래스)는 공개지만, 다른 auth manager 구현(FabAuthManager)은 공개가 아니에요.

auth managers와 나만의 것을 작성하는 방법에 대한 자세한 내용은 Auth manager에서 읽을 수 있어요.

Connections

Hooks를 만들 때 커스텀 Connections를 추가할 수 있어요. 커뮤니티 providers에서 구현된 사용 가능한 Connections는 Connections에서 확인할 수 있어요.

Hooks를 만들 때 tasks가 실행될 때 표시되는 커스텀 Extra Links를 추가할 수 있어요. 커뮤니티 providers에서 구현된 사용 가능한 extra links도 보여주는 Extra Links에서 자세한 내용을 찾을 수 있어요.

Logging and Monitoring

Airflow가 로그를 작성하는 방식을 확장할 수 있어요. 로그 작성에 대한 자세한 내용은 Logging & Monitoring에서 찾을 수 있어요.

커뮤니티 providers에서 구현된 사용 가능한 log writers도 보여주는 Writing logs.

Decorators

Dag 작성자는 TaskFlow 개념을 사용해 Dags를 작성하기 위해 decorators를 사용할 수 있어요. 모든 Decorators는 TaskDecorator에서 파생돼요.

Dag 작성자를 위한 기본 decorators는 이제 airflow.sdk 네임스페이스에 있어요: dag(), task(), asset(), setup(), task_group(), teardown(), chain(), chain_linear(), cross_downstream(), get_current_context()get_parsing_context().

Airflow에는 공개로 간주되는 Decorators 세트가 있어요. 그것들을 확장해 기능을 자유롭게 넓힐 수 있어요:

Note

Decorators는 이제 airflow.sdk 네임스페이스의 일부예요. 자세한 API 문서는 Task SDK 참조를 참고하세요.

커스텀 Decorators 만들기에 대한 자세한 내용은 Creating Custom @task Decorators에서 읽을 수 있어요.

Email notifications

Airflow에는 이메일 알림을 보내는 내장 방식이 있으며, 커스텀 이메일 알림 클래스를 추가해 확장할 수 있어요. 이메일 알림에 대한 자세한 내용은 Email Configuration에서 읽을 수 있어요.

Notifications

Airflow에는 다양한 on_*_callback을 사용해 알림을 보내는 내장 확장 가능한 방식이 있어요. 알림에 대한 자세한 내용은 Creating a Notifier에서 읽을 수 있어요.

Cluster Policies

Cluster Policies는 파싱되는 Dags나 실행되는 tasks에 클러스터 전반 정책을 동적으로 적용하는 방법이에요. Cluster Policies에 대한 자세한 내용은 Cluster Policies에서 읽을 수 있어요.

Lineage

Airflow는 데이터의 기원, 데이터에 일어나는 일, 시간 경과에 따라 어디로 이동하는지 추적하는 데 도움을 줄 수 있어요. lineage에 대한 자세한 내용은 Lineage에서 읽을 수 있어요.

Apache Airflow의 Public Interface가 아닌 것은 무엇인가요?

이 문서에서 언급되지 않은 모든 것은 non-Public Interface로 간주해야 해요.

때로는 다른 애플리케이션에서 그런 컴포넌트가 역호환성을 유지한다고 신뢰될 수 있지만, Airflow에서는 Public Interface의 일부가 아니며 언제든 변경될 수 있어요:

  • 데이터베이스 구조는 내부 구현 세부 사항으로 간주되며, 구조가 역호환 방식으로 유지될 것이라고 가정해서는 안 돼요.
  • Web UI는 계속 진화하며 HTML 요소에 대한 역호환 보장이 없어요.
  • 이 문서에서 명시적으로 언급된 것 외의 Python 클래스는 내부 구현 세부 사항으로 간주되며 역호환 방식으로 유지될 것이라고 가정해서는 안 돼요.

Dag 작성자가 작성한 코드의 메타데이터 데이터베이스 직접 접근은 더 이상 허용되지 않아요. Dag 작성자가 작성한 코드는 Dag 상태, task 히스토리, Dag runs를 조회하기 위해 메타데이터 데이터베이스에 직접 접근할 수 없어요 — 워커는 오로지 Execution API를 통해서만 통신해요. 대신 다음 대안 중 하나를 사용하세요:

  • Task Context: get_current_context()를 사용해 task instance 정보와 get_dr_count(), get_dagrun_state(), get_task_states() 같은 메서드에 접근해요.
  • REST API: Airflow 메타데이터에 프로그래밍 방식으로 접근하려면 Stable REST API를 사용하세요.
  • Python Client: Airflow와 Python 기반 상호작용에는 Python Client를 사용하세요.

이 변경은 아키텍처 분리를 개선하고 원격 실행 기능을 가능하게 해요.

직접 데이터베이스 접근 대신 Task Context를 사용하는 예시:

from airflow.sdk import dag, get_current_context, task, DagRunState
from datetime import datetime


@dag(dag_id="example_dag", start_date=datetime(2025, 1, 1), schedule="@hourly", tags=["misc"], catchup=False)
def example_dag():

    @task(task_id="check_dagrun_state")
    def check_state():
        context = get_current_context()
        ti = context["ti"]
        dag_run = context["dag_run"]

        # Use Task Context methods instead of direct DB access
        dr_count = ti.get_dr_count(dag_id="example_dag")
        dagrun_state = ti.get_dagrun_state(dag_id="example_dag", run_id=dag_run.run_id)

        return f"Dag run count: {dr_count}, current state: {dagrun_state}"

    check_state()


example_dag()

더 알아보기 (Learn more)