LLM WikiAccess-protected knowledge portal

WIKI

Airflow 3.3 멀티언어 태스크 SDK(AIP-108): Java·Go 실행 경계를 Python DAG 안에 들이는 법

왜 지금 봐야 하나 Airflow는 오케스트레이터지 실행기가 아니라는 말이 있다. 그러나 현실에서 Java 배치 로직이나 Go 바이너리를 DAG에 붙일 때는 BashOperator로 java jar 명령을 날리거나, PythonOperator 안에서 subprocess.run 을 호출하는 방식을 오래 써왔다. 이 구조는 단순하지만 문제가 있다. Airflow의 XCom, Variable, Connection 추상화가 서브프로세

경로human/study/content/database-frontier/56-airflow-3-3-language-sdk-java-go-coordinator-aip-108.md
카테고리Study
태그#aip #airflow #cicd #coordinator #java #mysql #sdk #study

왜 지금 봐야 하나

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 레이어에서 처리한다.


전체 흐름 개요

Python DAG
@task.stub(queue="java")
process_order(order_id)
▼ 태스크 인스턴스 생성
Airflow Scheduler
queue="java" 워커로
태스크 디스패치
Java/Go 워커 노드
JavaCoordinator
(JDK 17+ 필요)
Execution API 클라이언트
(XCom 읽기/쓰기 프록시)
@Task
process_order(Context ctx)
Airflow Core (Execution API)
XCom Store
Variable Store
Connection Store
Airflow 3.3 AIP-108 Coordinator 아키텍처

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 프록시 경계

Java/Go 태스크 내부
ctx.xcom().push("result", value)
▼ HTTP PUT /xcoms
Execution API
(워커→코어 사이드카)
▼ gRPC / DB 기록
XCom Store
(Airflow 메타데이터 DB)
다운스트림 Python 태스크
result = process_order.output
▼ XCom.pull()
Python dict로 역직렬화
Execution API XCom 라운드트립 흐름

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-schema

Docker 기반 배포라면 워커 이미지에 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_pushctx.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 및 와이어 프로토콜이 다음 릴리스에서 변경될 수 있다"고 명시되어 있다. 실 적용 전에 고려할 사항:

실 도입 전 단계로는 신규 서비스의 경량 태스크(파라미터 변환, 스키마 검증)에 먼저 적용하고, 기존 Java 배치 잡을 한꺼번에 마이그레이션하기보다는 점진적으로 교체하는 접근이 합리적이다.


Open Questions


References