Airflow 3.3 멀티언어 태스크 SDK(AIP-108): Java·Go 실행 경계를 Python DAG 안에 들이는 법
왜 지금 봐야 하나
Airflow는 오케스트레이터지 실행기가 아니라는 말이 있다. 그러나 현실에서 Java 배치 로직이나 Go 바이너리를 DAG에 붙일 때는 BashOperator로 java -jar 명령을 날리거나, PythonOperator 안에서 subprocess.run을 호출하는 방식을 오래 써왔다. 이 구조는 단순하지만 문제가 있다. Airflow의 XCom, Variable, Connection 추상화가 서브프로세스 경계에서 끊어진다. 파라미터는 환경 변수나 파일로 직렬화해야 하고, 실패 시 리트라이 시맨틱도 엉성해진다.
AIP-108은 이 간극을 Coordinator 프로세스와 Execution API 프록시로 채운다. 2026년 7월 6일 릴리스된 Airflow 3.3.0에 실험적 기능으로 포함되었다. Python DAG 개발자는 @task.stub(queue="java") 한 줄로 태스크를 외부 Coordinator에 위임한다. Java나 Go 개발자는 각자의 언어로 태스크 함수를 정의한다. 실행 경계는 Coordinator가 관리하고, XCom 입출력은 Execution API를 통해 Python 쪽과 동일한 방식으로 동작한다.
기존 방식의 한계
BashOperator / PythonOperator subprocess 패턴에서 반복되는 불편함을 정리하면 세 가지다.
| 문제 | 내용 |
|---|---|
| XCom 단절 | java -jar 호출 결과를 XCom으로 올리려면 stdout 파싱이나 임시 파일 경유 필요 |
| 실패 추적 어려움 | JVM 크래시나 Go 패닉이 Airflow 태스크 로그에 온전히 전달되지 않음 |
| 환경 의존성 관리 | 워커 머신에 특정 Java 버전이나 Go 바이너리가 설치되어 있어야 함 |
| 타입 정보 유실 | Python dict → JSON 문자열 → 환경 변수 → Java String → 다시 파싱하는 직렬화 단계 |
AIP-108은 Coordinator 프로세스를 사이에 두어 이 네 가지를 Execution API 레이어에서 처리한다.
전체 흐름 개요
process_order(order_id)
태스크 디스패치
(JDK 17+ 필요)
(XCom 읽기/쓰기 프록시)
process_order(Context ctx)
DAG 정의 시점에는 @task.stub만 있으면 된다. 실제 비즈니스 로직은 Java(또는 Go) SDK 모듈에 구현된다. 스케줄러가 태스크 인스턴스를 queue="java" 워커에 배치하면, 그 워커에서 JavaCoordinator가 JVM 서브프로세스를 스폰하여 태스크를 실행한다.
Python 쪽: @task.stub
DAG 파일에서 Java 태스크를 선언하는 최소 코드다.
from airflow.sdk import dag, task
@dag
def order_pipeline():
@task.stub(queue="java")
def process_order(order_id: str) -> dict:
... # 본체는 Java에서 구현
@task.stub(queue="go")
def validate_schema(payload: dict) -> bool:
... # 본체는 Go에서 구현
result = process_order(order_id="{{ dag_run.conf['order_id'] }}")
validate_schema(result)@task.stub은 DAG 파싱 시 해당 태스크의 의존성 그래프만 등록한다. 함수 본체(...)는 무시된다. queue 파라미터가 Airflow에게 "이 태스크는 외부 Coordinator가 처리한다"는 신호를 준다.
XCom은 Python 태스크와 완전히 호환된다. process_order 반환 dict가 XCom에 저장되면, validate_schema가 그것을 Python dict로 받는다. Java→Python XCom 왕복에서 직렬화는 Execution API가 처리한다.
Java 쪽: JavaCoordinator
의존성 추가
Maven 프로젝트 기준:
<dependency>
<groupId>org.apache.airflow</groupId>
<artifactId>airflow-sdk</artifactId>
<version>3.3.0</version>
</dependency>Gradle:
implementation 'org.apache.airflow:airflow-sdk:3.3.0'JDK 17 이상이 필요하다. 내부적으로 Java 17 record, sealed class를 사용한다.
태스크 구현
import org.apache.airflow.sdk.Task;
import org.apache.airflow.sdk.Context;
import org.apache.airflow.sdk.JavaCoordinator;
import java.util.Map;
public class OrderTasks {
@Task(name = "process_order")
public Map<String, Object> processOrder(Context ctx) {
String orderId = (String) ctx.xcom().pull("order_id");
// 비즈니스 로직
return Map.of("status", "processed", "order_id", orderId);
}
public static void main(String[] args) throws Exception {
JavaCoordinator.run(OrderTasks.class);
}
}JavaCoordinator.run()은 워커가 태스크 인스턴스를 전달할 때까지 대기한다. 태스크 디스패치를 받으면 @Task(name=...) 어노테이션으로 매핑된 메서드를 실행한다. Context는 Execution API를 통해 XCom, Variable, Connection에 접근하는 핸들이다.
Go 쪽: ExecutableCoordinator
의존성 추가
go get github.com/apache/[email protected]태스크 구현
package main
import (
"context"
"github.com/apache/airflow-sdk-go/airflow"
)
func validateSchema(ctx context.Context, ac airflow.Context) (bool, error) {
payload, err := ac.XCom().Pull(ctx, "payload")
if err != nil {
return false, err
}
// 스키마 검증 로직
_ = payload
return true, nil
}
func main() {
airflow.Register("validate_schema", validateSchema)
airflow.Run(context.Background())
}Go의 경우 Coordinator는 외부 프로세스가 아닌 airflow.Run()을 호출하는 단일 바이너리다. airflow.Register로 태스크 이름과 함수를 매핑하면 ExecutableCoordinator가 워커 프로세스와 소켓 통신을 처리한다.
Execution API: XCom 프록시 경계
(워커→코어 사이드카)
(Airflow 메타데이터 DB)
Coordinator와 Airflow 코어 사이에는 Execution API가 있다. 이 API는 태스크 실행 중 사이드카 HTTP 서버로 동작하며, XCom 읽기/쓰기, Variable 조회, Connection 자격 증명 조회를 프록시한다. Java나 Go 코드는 직접 메타데이터 DB에 접근하지 않는다.
이 구조는 Airflow 3.x에서 도입된 Task Execution Interface(AIP-72) 위에 구축된다. AIP-108은 "외부 언어로 작성된 태스크도 이 인터페이스를 통해 Airflow 컨텍스트를 사용할 수 있어야 한다"는 확장이다.
워커 설정
기존 Celery나 LocalExecutor 워커와 동일한 노드에서 Coordinator를 실행한다. Airflow 워커가 queue="java" 태스크를 수신하면, 미리 등록된 Coordinator 바이너리(또는 JAR)를 서브프로세스로 실행한다.
airflow.cfg 또는 환경 변수로 Coordinator 바이너리 경로를 지정한다:
[language_sdk]
java_coordinator_jar = /opt/airflow/coordinators/order-tasks-3.3.0.jar
go_coordinator_bin = /opt/airflow/coordinators/validate-schemaDocker 기반 배포라면 워커 이미지에 JDK 17과 Coordinator JAR을 함께 패키징한다. Go 바이너리는 정적 컴파일이므로 이미지 크기 영향이 작다.
BashOperator 마이그레이션 패턴
이전 패턴과 AIP-108 이후를 비교한다.
| 구분 | 이전 (BashOperator) | 이후 (AIP-108) |
|---|---|---|
| DAG 선언 | BashOperator(bash_command="java -jar ... --order-id {{ ... }}") | @task.stub(queue="java") |
| 입력 전달 | 환경 변수, 커맨드라인 인수 | ctx.xcom().pull() |
| 결과 반환 | stdout 파싱 후 xcom_push | ctx.xcom().push() |
| 오류 전파 | exit code → Airflow 실패 | Exception → Airflow 실패 + 스택트레이스 |
| 재시도 | Airflow 재시도 = BashOperator 재실행 | Airflow 재시도 = Coordinator 재시작 |
| 타입 안전성 | 없음 | Java record / Go struct → JSON 직렬화 |
BashOperator로 java -jar를 호출하던 태스크를 마이그레이션할 때, 진입점을 JavaCoordinator.run()으로 바꾸고 환경 변수 접근을 ctx.xcom().pull() / ctx.variable().get()으로 치환하는 것이 핵심이다.
실험적 기능임을 감안한 운영 판단
AIP-108은 Airflow 3.3.0에서 experimental로 표시된다. 공식 문서에 "API 및 와이어 프로토콜이 다음 릴리스에서 변경될 수 있다"고 명시되어 있다. 실 적용 전에 고려할 사항:
- Execution API 안정성: AIP-72 기반 Execution API 자체는 Airflow 3.x에서 안정화 진행 중이다. AIP-108은 그 위에 얹힌 언어 바인딩이어서 API가 먼저 굳어야 SDK도 안정화된다.
- 관찰 가능성: Coordinator 서브프로세스 로그가 Airflow 태스크 로그에 통합되지만, 상세 수준은 BashOperator보다 좋다. Java 스택트레이스 전체가 태스크 로그에 나타난다.
- 롤백 경로:
@task.stub을BashOperator로 되돌리는 것은 DAG 파일 변경 한 줄이다. Coordinator JAR은 별도 배포이므로 워커 이미지 롤백도 독립적이다. - Kubernetes 워커: KubernetesExecutor에서는
queue="java"태스크를 처리하는 전용 파드 스펙을pod_template_file로 지정하면 된다. 파드에 JDK와 JAR을 마운트한다.
실 도입 전 단계로는 신규 서비스의 경량 태스크(파라미터 변환, 스키마 검증)에 먼저 적용하고, 기존 Java 배치 잡을 한꺼번에 마이그레이션하기보다는 점진적으로 교체하는 접근이 합리적이다.
Open Questions
- Execution API 와이어 프로토콜 확정 일정 (현재 experimental): Airflow 3.4 또는 4.0에서 안정화 예상이나 미확인
- JavaCoordinator 멀티스레딩 지원 여부: 단일 태스크 인스턴스당 하나의 JVM인지, 하나의 JVM에서 여러 태스크를 처리하는지 공식 문서에 불명확
- Go SDK의
context.Context취소 전파가 Airflow 태스크 타임아웃과 정확히 연동되는지 검증 필요 - Python 3.13+ 호환성: Airflow 3.3.0 릴리스 노트에 Python 3.13 지원 명시 여부 확인 필요
References
- Apache Airflow 3.3.0 릴리스 노트: https://airflow.apache.org/docs/apache-airflow/3.3.0/release_notes.html
- AIP-108 Language Task SDK 디자인 문서: https://github.com/apache/airflow/blob/main/AIP-108.md
- Apache Airflow Java SDK (org.apache.airflow:airflow-sdk): https://github.com/apache/airflow-java-sdk
- Apache Airflow Go SDK: https://github.com/apache/airflow-sdk-go
- AIP-72 Task Execution Interface: https://github.com/apache/airflow/blob/main/AIP-72.md
- Airflow 3.x Execution API 개요: https://airflow.apache.org/docs/apache-airflow/stable/execution-api.html
- airflow.cfg language_sdk 설정 레퍼런스: https://airflow.apache.org/docs/apache-airflow/3.3.0/configurations-ref.html
- Airflow 커뮤니티 디스커션 (AIP-108 initial proposal): https://github.com/apache/airflow/discussions/40000
- Airflow Slack #aip-108 채널 요약 (Apache Airflow Slack): https://apache-airflow.slack.com
- Airflow 3.3 마이그레이션 가이드: https://airflow.apache.org/docs/apache-airflow/3.3.0/migration-guide.html