비-Python Task SDK
비-Python Task SDK (Non-Python Task SDKs)
이 페이지는 개별 Task 구현을 Python이 아닌 다른 언어로 작성할 수 있게 해주는 비-Python Task SDK(실험적 기능)를 다뤄요. DAG는 항상 Python으로 정의되지만, Task 로직은 다른 언어의 런타임이 실행해요. Stub Task, Coordinator(Java·Go), Coordinator 설정, 새 컴파일 언어 SDK 구현 방법을 설명해요.
출처: 문서
본문
이것은 실험적 기능이에요.
Airflow DAG는 항상 Python으로 정의되지만, 개별 Task 구현은 다른 언어로 작성할 수 있어요. Task가 실행될 때 Airflow의 worker가 target 언어를 호출해 Task 로직을 실행하고, 그 결과를 다시 통신으로 전달해요. Dag 작성자는 Task가 어디에 있는지 선언하기 위해 가벼운 Python stub을 사용해요. 실제 비즈니스 로직, Airflow API 호출, 라이브러리 의존성을 포함한 나머지 모든 것은 비-Python 구현에 있어요.
사용 가능한 언어 SDK
| 언어 | Coordinator 클래스 | 최소 런타임 | 가이드 |
|---|---|---|---|
| JVM 언어 (예: Java) | airflow.sdk.coordinators.java.JavaCoordinator |
JRE 17 | Java SDK |
| Go | airflow.sdk.coordinators.executable.ExecutableCoordinator |
없음 (네이티브 바이너리) | Go SDK |
작동 원리
실행 모델에는 세 가지 움직이는 부분이 있어요.
Dag의 Stub Task
: Dag 파일은 @task.stub를 사용해 Task를 선언해요. stub은 scheduler 관점에서 보면 일반 Airflow Task예요. 의존성, 재시도, pool, 그리고 다른 모든 Task 레벨 기능에 다른 @task 데코레이트 Python 함수와 정확히 동일하게 참여해요. 유일한 차이는 worker가 함수 정의 안의 Python 코드를 실행하지 않고, 대신 coordinator에게 실행을 위임한다는 점이에요.
Coordinators
: Coordinator는 [sdk] coordinators 설정에 등록된 Python 객체예요. 이는 Airflow worker의 일부로 간주돼요. Worker가 stub task를 집어 들면, 그 Task의 지정된 queue에 매핑된 coordinator를 찾아 coordinator를 사용해 Task를 실행해요. Coordinator는 target 언어의 런타임을 관리하고, 메시지를 전달하며, 결과를 Airflow로 다시 중계할 책임이 있어요. 모든 coordinator는 task-sdk:airflow.sdk.execution_time.coordinator.BaseCoordinator를 확장해요.
언어 런타임 : Coordinator는 task instance당 하나의 수명이 짧은 런타임을 호출해요. 대부분의 경우 이것은 비-Python 언어로 구현된 실행 파일의 서브프로세스예요. 런타임은 worker에서 메시지를 받아 워크로드를 식별하고, Task를 실행하며, proxy 역할을 하는 coordinator를 통해 worker 프로세스로 다시 통신해요.
Stub Task
Stub task는 @task.stub 데코레이터로 선언돼요. 여전히 Python task 선언이므로, 일반 Dag·Task에서 사용 가능한 모든 파라미터가 적용돼요. Task 의존성도 Python Dag 파일에서 정의돼요. Scheduler는 stub을 다른 task처럼 취급해요.
import datetime
from airflow.sdk import dag, task
@dag
def my_pipeline():
raw = fetch_data() # normal Python task
@task.stub(
queue="java", # routes to the JavaCoordinator
retries=3,
retry_delay=datetime.timedelta(minutes=5),
execution_timeout=datetime.timedelta(hours=1),
pool="heavy_tasks",
)
def process(raw_value): ... # implemented in Java
@task.stub(queue="java")
def export(processed_value): ...
export(process(raw))
my_pipeline()
queue 파라미터가 어느 coordinator가 Task를 처리할지 결정해요. 다른 모든 @task 키워드 인자는 task instance에 저장되고 평소처럼 Airflow의 scheduler와 worker가 존중해요.
Stub task가 생성한 XCom 값은 다운스트림 Python task에서 볼 수 있고 그 반대도 마찬가지예요. 다만, XCom 참조는 Python Dag 안에 정의되어야 하지만(그것들은 task 의존성입니다), 언어 구현에서 실제로 값을 읽어야 하며 그 반대도 마찬가지예요. 올바르게 하는 방법은 특정 언어 SDK 문서를 참고해요.
Coordinator 설정
Coordinator는 [sdk] 아래 airflow.cfg(또는 환경 변수)에 등록돼요.
coordinators
: 논리적 coordinator 이름을 그 클래스와 키워드 인자에 매핑하는 JSON 객체:
````ini
[sdk]
coordinators = {
"my-coordinator": {
"classpath": "path.to.CoordinatorClass",
"kwargs": {},
"extra": {}
}
}
````
`classpath` 값은 worker가 import할 수 있어야 해요. `kwargs`는 coordinator의 생성자에 직접 전달돼요. 각 coordinator의 허용 kwargs는 언어별 가이드를 참고해요 (예: [`JavaCoordinator`](/docs/task-sdk/stable/api.html#airflow.sdk.coordinators.java.JavaCoordinator "(in Apache Airflow Task SDK v1.4.0)")에 대한 [JavaCoordinator 설정](java.html#java-sdk-coordinator-config)).
`extra`는 coordinator 인스턴스에 결합하지 않고 coordinator와 연관시키고 싶은 추가 정보를 위한 선택적 객체예요. Coordinator 자체는 그것을 절대 받지 않아요. 다른 컴포넌트가 필요에 따라 읽어요. 예를 들어 KubernetesExecutor는 `extra.pod_template_file`을 읽어 특정 pod 템플릿에서 queue의 worker pod를 시작하고, `extra.worker_container_repository` + `extra.worker_container_tag`를 읽어 그 queue의 worker 기본 이미지를 덮어써요 (두 키 모두 필요) — 예를 들어 Java coordinator를 위한 JVM을 번들한 이미지요.
queue_to_coordinator
: Celery queue 이름을 coordinator 이름에 매핑하는 JSON 객체:
````ini
[sdk]
queue_to_coordinator = {"jdk17": "my-coordinator"}
````
Stub에 `queue="jdk17"`이 있는 Task는 `"my-coordinator"`라는 coordinator로 디스패치돼요. 하나의 coordinator가 여러 queue를 제공할 수 있고, 하나의 queue는 하나의 coordinator에만 매핑될 수 있어요.
두 설정 모두 표준 Airflow 규칙을 사용해 환경 변수로 제공할 수 있어요:
AIRFLOW__SDK__COORDINATORS='{"my-coordinator": {...}}'
AIRFLOW__SDK__QUEUE_TO_COORDINATOR='{"jdk17": "my-coordinator"}'
새 컴파일 언어 SDK 구현하기
airflow.sdk.coordinators.executable.ExecutableCoordinator는 번들 파일을 직접 실행해 Task를 실행해요. 따라서 빌드 산출물이 worker가 추가 런타임 의존성 없이 실행할 수 있는 독립 실행형 바이너리인 컴파일 언어에만 맞아요 — 즉 바이너리 실행에 언어 런타임, 가상 머신, 인터프리터가 worker에 설치될 필요가 없어야 해요 (Go, Rust, C, C++, Zig, …). 산출물이 실행 시점에 런타임이 여전히 필요한 언어는 이 coordinator에 맞지 않아요. JVM 언어는 예를 들어 JRE가 필요한 바이트코드로 컴파일되므로, airflow.sdk.coordinators.java.JavaCoordinator가 대신 지원해요.
새 그런 언어를 지원하려면 coordinator가 소비하는 공유 온디스크 형식으로 bundle을 만들고 coordinator IPC 프로토콜(--comm / --logs 소켓 인자)을 사용해야 해요. 실행 파일에 추가되는 AFBNDL01 footer, 바이너리 무결성 해시, dag_id·task_id의 airflow-metadata.yaml 매니페스트를 포함한 그 형식은, 읽기 알고리즘과 호환성/버전 관리 규칙과 함께 Executable Bundle Spec에 명시돼 있어요. 그 페이지는 또한 빌드 도구와 검증기가 사용할 매니페스트의 머신 판독 가능 JSON Schema도 게시해요. 스펙을 따르면 scheduler, worker, UI 변경 없이 새 언어의 번들을 Airflow가 발견할 수 있어요. Go SDK가 참조 구현으로 사용되는 워크드 에세이에요.