Java SDK
Java SDK
이 페이지는 Airflow Task 로직을 Java·Kotlin·기타 JVM 언어로 구현할 수 있게 해주는 Java SDK(실험적 기능)를 다뤄요. DAG와 스케줄링은 Python에 남고, 개별 Task는 각 task instance마다 JavaCoordinator가 실행하는 JVM 서브프로세스에 위임해요. Task 작성(애노테이션 기반/인터페이스 기반 API), 로깅, XCom 타입 매핑, 빌드·패키징(Gradle/Maven)을 설명해요.
출처: 문서
본문
이것은 실험적 기능이에요.
Java SDK는 Airflow Task 로직을 Java, Kotlin, 또는 기타 JVM 언어로 구현할 수 있게 해줘요. DAG와 스케줄링은 Python에 남고, 개별 Task는 각 task instance마다 JavaCoordinator가 생성하는 JVM 서브프로세스에 위임해요.
API 참조
Java SDK에 대한 생성된 API 참조는 Airflow 문서와 함께 Java SDK API Reference에 게시돼 있어요.
선행 조건
- Airflow worker 노드에 JRE 17 이상이 있어야 해요.
- 컴파일된 Task JAR(들)과 JVM 의존성이 worker에서 접근 가능해야 해요.
apache-airflow-task-sdk패키지(Airflow와 함께 설치)가 coordinator를 제공해요. 추가 Python 패키지는 필요 없어요.
빠른 시작
다음 예제는 최소한의 움직이는 부분을 보여줘요: 두 개의 stub task를 가진 Python Dag와 그 Task들의 Java 구현.
Python DAG (스케줄링 쪽)
from airflow.sdk import dag, task
@dag
def sales_pipeline():
@task.stub(queue="java")
def extract(): ...
@task.stub(queue="java")
def transform(extracted): ...
@task()
def load(transformed):
print(f"Loaded: {transformed}")
load(transform(extract()))
sales_pipeline()
Java 구현
import org.apache.airflow.sdk.*;
@Builder.Dag(id = "sales_pipeline")
public class SalesPipeline {
@Builder.Task(id = "extract")
public long extract(Client client) {
var conn = client.getConnection("sales_db");
// ... fetch data using conn.host, conn.login, conn.password ...
return recordCount;
}
@Builder.Task(id = "transform")
public long transform(
Client client,
@Builder.XCom(task = "extract") long recordCount
) {
var threshold = (String) client.getVariable("transform_threshold");
// ... process data ...
return transformedCount;
}
}
Note
Python과 Java 모두에서
transform이 상위 XCom을 받는 인자를 가져야 하는지 주목해요. Python 쪽은 의존성을 선언하는 데 필요하고, Java 쪽은 실제로 값을 검색하는 데 필요해요.
Java 진입점
public class Main implements BundleBuilder {
@Override
public Iterable<Dag> getDags() {
return List.of(SalesPipelineBuilder.build()); // SalesPipelineBuilder generated at compile time
}
public static void main(String[] args) {
Server.create(args).serve(new Main().build());
}
}
Coordinator 구성
[sdk]
coordinators = {
"java-jdk17": {
"classpath": "airflow.sdk.coordinators.java.JavaCoordinator",
"kwargs": {"jars_root": ["/opt/airflow/jars"]}
}
}
queue_to_coordinator = {"java": "java-jdk17"}
허용된 kwargs 전체 목록은 JavaCoordinator 구성을 참고해요.
Task 작성
Java SDK는 Task를 구현하는 두 가지 API를 제공해요. 둘 다 같은 런타임 동작을 만들며, 선택은 스타일 문제예요.
애노테이션 기반 API
일반 Java 클래스에 애노테이션을 붙이고 SDK가 컴파일 시점에 보일러플레이트를 생성하게 해요.
| 애노테이션 | 용도 |
|---|---|
@Builder.Dag(id = "...") |
클래스를 task 컨테이너로 표시해요. id는 Python DAG의 dag_id와 일치해야 해요. |
@Builder.Task(id = "...") |
메서드를 task 구현으로 표시해요. id는 Python DAG의 @task.stub 함수 이름과 일치해야 해요. id를 생략하면 메서드 이름이 사용돼요. |
@Builder.XCom(task = "...") |
이름 있는 상위 Task의 return_value XCom을 메서드 파라미터로 주입해요. 파라미터 타입은 저장된 값과 호환되어야 해요 (XCom 타입 매핑 참고). |
애노테이션 프로세서가 task 레지스트리를 연결하고 XCom 주입을 자동 처리하는 <ClassName>Builder 클래스를 생성해요.
@Builder.Dag(id = "my_dag")
public class MyDag {
@Builder.Task(id = "fetch")
public String fetch(Client client) throws Exception {
var conn = client.getConnection("my_api");
// implement task logic
return result;
}
@Builder.Task(id = "process")
public long process(
Client client,
@Builder.XCom(task = "fetch") String fetched
) {
var threshold = (String) client.getVariable("process_threshold");
// implement task logic
return count;
}
}
Task 메서드는 throws Exception을 선언할 수 있어요. 잡히지 않은 예외는 Airflow에서 task instance를 실패로 표시해요(stub에 재시도가 구성되어 있으면 재시도 트리거).
인터페이스 기반 API
Task가 등록되는 방식과 XCom이 읽히는 방식을 완전히 제어하려면 Task 인터페이스를 직접 구현해요.
import org.apache.airflow.sdk.*;
public class FetchTask implements Task {
@Override
public void execute(Context context, Client client) throws Exception {
var conn = client.getConnection("my_api");
// implement task logic
client.setXCom(result);
}
}
BundleBuilder에서 Task를 수동으로 등록해요:
public class MyBundle implements BundleBuilder {
@Override
public Iterable<Dag> getDags() {
var dag = new Dag("my_dag");
dag.addTask("fetch", FetchTask.class);
dag.addTask("process", ProcessTask.class);
return List.of(dag);
}
}
자세한 내용은 Java SDK API Reference를 참고해요.
로깅
Task 코드는 어떤 일반적인 Java 로깅 프레임워크로도 로그 레코드를 내보낼 수 있어요. SDK는 그 레코드를 Airflow의 task log 저장소로 전달하는 선택적 통합 라이브러리를 제공하며, 표준 task 출력과 함께 Airflow UI에 나타나요.
Task 클래스의 정적 필드로 클래스의 자체 타입을 이름으로 사용해 로거를 선언해요. 어떤 로깅 프레임워크를 선택하든 관례적인 패턴이에요:
private static final System.Logger log =
System.getLogger(SalesPipeline.class.getName());
@Builder.Task(id = "extract")
public long extract(Client client) {
log.log(System.Logger.Level.INFO, "Starting extraction");
return recordCount;
}
아래 Gradle 스니펫은 의존성 선언을 보여줘요. 모든 Airflow 아티팩트 버전은 airflow-sdk-bom이 관리해요. Maven 사용자는 Maven의 패턴을 따라 같은 아티팩트 ID를 적용해요.
System.Logger (Java Platform Logging)
Java 9의 새 로깅 파사드 java.lang.System.Logger(JEP 264, 흔히 JPL로 약칭)는 서드파티 API를 끌어들이지 않고 라이브러리가 사용할 수 있어요. airflow-sdk-jpl 아티팩트는 ServiceLoader로 AirflowSystemLoggerFinder를 등록해, 모든 System.Logger 호출을 Airflow의 task log 저장소로 직접 라우팅해요.
implementation("org.apache.airflow:airflow-sdk-jpl:${version}")
설정 파일이나 시작 호출이 필요 없어요. JAR이 클래스패스에 있는 한 ServiceLoader 매커니즘이 provider를 자동으로 발견해요.
Note
airflow-sdk-jpl과 함께 두 번째System.LoggerFinder구현을 추가하지 마세요. JVM은ServiceLoader로 하나의 finder를 선택하며, 클래스패스에 여러 provider가 있으면 예측 불가능한 동작이 발생해요.
SLF4J 2.x
SLF4J 바인딩은 ServiceLoader로 자동 발견돼요. 설정 파일이나 시작 호출이 필요 없어요.
implementation("org.apache.airflow:airflow-sdk-slf4j:${version}")
위는 자동으로 SLF4J API를 끌어들이므로 slf4j-api를 직접 추가할 필요가 없어요.
Note
airflow-sdk-slf4j와 함께 두 번째 SLF4J 바인딩(예:logback-classic이나slf4j-simple)을 추가하지 마세요. SLF4J 2.x는 여러 바인딩에 대해 경고하고 가장 예측 불가능하게 하나를 선택해요.
Log4j 2
airflow-sdk-log4j2는 log4j-api를 전이적 의존성으로 선언하므로 후자를 별도로 추가할 필요가 없어요. 또한 시작 시 airflow-sdk-log4j2가 제공하는 커스텀 AirflowAppender를 발견하는 플러그인 로더를 호스팅하려면 log4j-core를 런타임 클래스패스에 둬야 해요:
implementation("org.apache.airflow:airflow-sdk-log4j2:${version}")
runtimeOnly("org.apache.logging.log4j:log4j-core:${log4jVersion}")
log4j2.xml에 AirflowAppender를 선언해요:
<?xml version="1.0" encoding="UTF-8"?>
<Configuration>
<Appenders>
<AirflowAppender name="Airflow"/>
</Appenders>
<Loggers>
<Root level="info">
<AppenderRef ref="Airflow"/>
</Root>
</Loggers>
</Configuration>
java.util.logging
아티팩트를 추가해요:
implementation("org.apache.airflow:airflow-sdk-jul:${version}")
그리고 Task가 실행되기 전에 시작 시 AirflowJulHandler.setup()을 호출해요. 이는 JUL root logger의 기존 핸들러(기본 ConsoleHandler 포함 — 그 stderr 출력을 Airflow가 ERROR 레벨에서 task.stderr로 캡처해 각 레코드를 중복하고 레벨을 잘못 표시할 것)를 지우고 그 자리에 AirflowJulHandler를 설치해요:
public static void main(String[] args) {
AirflowJulHandler.setup();
Server.create(args).serve(new MyBundle());
}
또는 logging.properties 파일에 핸들러를 선언하고 java.util.logging.config.file 시스템 속성(coordinator 설정의 jvm_args로 설정)으로 JUL을 그 파일에 가리켜요:
handlers = org.apache.airflow.sdk.jul.AirflowJulHandler
[sdk]
coordinators = {
"java-jdk17": {
"classpath": "airflow.sdk.coordinators.java.JavaCoordinator",
"kwargs": {
"jars_root": ["/opt/airflow/jars"],
"jvm_args": ["-Djava.util.logging.config.file=/opt/airflow/logging.properties"]
}
}
}
다른 프레임워크
몇 가지 일반적으로 사용되는 로깅 API는 전용 Airflow 아티팩트 없이도 다뤄져요:
- Logback은 그 자체가 SLF4J 바인딩이에요.
logback-classic을airflow-sdk-slf4j로 교체하면 Task 코드 변경이 필요 없어요. - **Apache Commons Logging (JCL)**은
org.slf4j:jcl-over-slf4j로 SLF4J에, 또는org.apache.logging.log4j:log4j-jcl로 Log4j 2에 브리지할 수 있어요.
XCom 타입 매핑
XCom 값은 Airflow 메타데이터 데이터베이스에 JSON으로 저장돼요. 아래 표는 getXCom으로 다시 읽을 때 JSON 타입이 Java 객체로 어떻게 표현되는지 보여줘요.
| Python 타입 | JSON | Java 타입 (getXCom으로부터) |
|---|---|---|
int |
number (integer) | Long (들어맞는 값의 경우; 그 외 BigInteger) |
float |
number (decimal) | Double |
str |
string | String |
bool |
boolean | Boolean |
None |
null | null |
list |
array | List<Object> |
dict |
object | Map<String, Object> |
빌드·패키징
Java SDK는 JAR로 배포돼요. 아래 섹션들은 Gradle이나 Maven으로 번들을 빌드하는 방법을 보여줘요.
Gradle
build.gradle에 Airflow SDK Gradle 플러그인을 적용해요:
plugins {
id("org.apache.airflow.sdk") version "${version}"
}
dependencies {
annotationProcessor("org.apache.airflow:airflow-sdk-processor:${version}")
implementation("org.apache.airflow:airflow-sdk:${version}")
}
airflowBundle {
mainClass = "com.example.Main" // Point to your main class instead.
}
그런 다음 실행해요:
./gradlew bundle
build/bundle/ 디렉터리에 필요한 모든 JAR(들)이 있어요. 그것을 coordinator 설정의 jars_root가 가리키는 디렉터리로 복사하거나 마운트해요. JavaCoordinator는 jars_root를 재귀적으로 스캔하고 클래스패스를 자동으로 빌드해요.
Note
annotationProcessor항목은 애노테이션 기반 API를 사용할 때만 필요해요. 인터페이스 기반 API에는 필요하지 않아요.
Note
플러그인은 기본적으로 Shadow 플러그인으로 fat JAR을 생성해요. 프로젝트 간 의존성 문제를 피하기 위해 JAR 하나만 배포하므로 일반적으로 좋은 생각이에요. 이것이 맞지 않으면
airflowBundle에서fatJar = false를 설정해 대신 얇은(thin) JAR을 만들어요. 나머지 과정은 같지만, Airflow가jars_root로 찾을 수 있는 곳에 모든 의존성 JAR을 둬야 해요.
Maven
아티팩트 버전과 ${airflow.supervisor.schema.version} 속성이 한 곳에서 관리되도록 airflow-sdk-bom Bill of Materials를 import해요:
<dependencyManagement>
<dependencies>
<dependency>
<groupId>org.apache.airflow</groupId>
<artifactId>airflow-sdk-bom</artifactId>
<version>${version}</version>
<type>pom</type>
<scope>import</scope>
</dependency>
</dependencies>
</dependencyManagement>
SDK를 의존성으로 추가해요(버전은 BOM이 관리):
<dependencies>
<dependency>
<groupId>org.apache.airflow</groupId>
<artifactId>airflow-sdk</artifactId>
</dependency>
</dependencies>
애노테이션 프로세서를 maven-compiler-plugin을 통해 연결해 런타임 클래스패스에서 벗어나게 해요:
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-compiler-plugin</artifactId>
<configuration>
<annotationProcessorPaths>
<path>
<groupId>org.apache.airflow</groupId>
<artifactId>airflow-sdk-processor</artifactId>
<version>${version}</version>
</path>
</annotationProcessorPaths>
</configuration>
</plugin>
옵션 1 (권장): fat JAR
maven-shade-plugin을 사용해 코드와 모든 의존성을 단일 JAR로 번들해요. 이것이 가장 단순한 배포예요: 파일 하나, 런타임 의존성 관리 없음.
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-shade-plugin</artifactId>
<version>3.6.0</version>
<executions>
<execution>
<phase>package</phase>
<goals><goal>shade</goal></goals>
<configuration>
<transformers>
<transformer implementation="org.apache.maven.plugins.shade.resource.ManifestResourceTransformer">
<!-- Replace with your BundleBuilder implementation. -->
<mainClass>com.example.Main</mainClass>
<manifestEntries>
<!-- Resolved from the BOM; do not hard-code this value. -->
<Airflow-Supervisor-Schema-Version>${airflow.supervisor.schema.version}</Airflow-Supervisor-Schema-Version>
</manifestEntries>
</transformer>
</transformers>
</configuration>
</execution>
</executions>
</plugin>
그런 다음 실행해요:
mvn package
fat JAR은 target/<artifactId>-<version>.jar에 기록돼요. coordinator에서 jars_root로 구성된 디렉터리에 복사해요.
옵션 2: 별도 의존성을 가진 얇은 JAR
fat JAR이 프로젝트에 맞지 않으면 maven-jar-plugin으로 일반 JAR에 Main-Class를 설정하고 maven-dependency-plugin으로 모든 런타임 의존성을 그 옆에 수집해요. 여기서는 Airflow-Supervisor-Schema-Version을 설정할 필요가 없다는 점을 주의하세요. Airflow가 클래스패스의 airflow-sdk JAR에서 직접 읽으니까요.
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-jar-plugin</artifactId>
<configuration>
<archive>
<manifestEntries>
<Main-Class>com.example.Main</Main-Class>
</manifestEntries>
</archive>
</configuration>
</plugin>
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-dependency-plugin</artifactId>
<executions>
<execution>
<id>copy-dependencies</id>
<phase>package</phase>
<goals><goal>copy-dependencies</goal></goals>
<configuration>
<outputDirectory>${project.build.directory}/bundle</outputDirectory>
<includeScope>runtime</includeScope>
</configuration>
</execution>
<execution>
<id>copy-artifact</id>
<phase>package</phase>
<goals><goal>copy</goal></goals>
<configuration>
<artifactItems>
<artifactItem>
<groupId>${project.groupId}</groupId>
<artifactId>${project.artifactId}</artifactId>
<version>${project.version}</version>
<outputDirectory>${project.build.directory}/bundle</outputDirectory>
</artifactItem>
</artifactItems>
</configuration>
</execution>
</executions>
</plugin>
그런 다음 실행해요:
mvn package
target/bundle/에 얇은 JAR과 모든 런타임 의존성 JAR이 있을 거예요. jars_root를 이 디렉터리로 가리켜요.
Note
annotationProcessorPaths항목은 애노테이션 기반 API를 사용할 때만 필요해요.
Note
Gradle 플러그인과 달리 Maven에는
verifyBundleMainClass검증 단계와 동등한 것이 없어요. 잘못된<mainClass>값은 런타임까지 잡히지 않아요.
JavaCoordinator 구성
coordinators 설정 항목의 모든 kwargs는 JavaCoordinator 생성자로 전달돼요:
| 파라미터 | 기본값 | 설명 |
|---|---|---|
jars_root |
(필수) | .jar 파일을 재귀적으로 스캔하는 하나 이상의 디렉터리. 문자열, 경로, 문지열/경로 리스트를 허용해요. |
java_executable |
"java" |
java 바이너리 경로. $PATH의 java로 기본 설정돼요. |
jvm_args |
[] |
["-Xmx1g", "-Dsome.property=value"] 같은 추가 JVM 인자. |
main_class |
(자동 감지) | 명시적 진입점 클래스. 생략하면 JavaCoordinator가 jars_root에서 매니페스트가 Main-Class를 설정한 JAR을 스캔해요. 실행 가능한 JAR이 여러 개면 결과가 비결정적이에요. 그 경우 main_class를 명시적으로 설정하세요. |
task_startup_timeout |
10.0 |
실행 후 JVM 서브프로세스가 연결될 때까지 기다리는 초. JVM 시작이 느리면(예: 제한된 하드웨어나 큰 클래스패스) 이 값을 늘려요. |
Note
[sdk]설정은 시작 시 읽히므로,coordinators나queue_to_coordinator에 대한 변경(예:jvm_args추가)은 scheduler(또는airflow standalone)를 재시작한 후에만 적용돼요. 반대로 재빌드된 번들 JAR은 task instance마다 새 JVM이 실행되므로 재시작 없이 다음 task 실행에서 반영돼요.
Java 실행 파일 고정
일반적인 권장으로, $PATH에서 java가 해석되도록 의존하기보다 java_executable을 절대 경로로 설정해요. 이렇게 하면 Task가 알려진 JDK에 고정되는데, 이는 Airflow 관리자가 시스템 전역 java를 통제하지 못할 수 있는 production이나 기업 환경에서 가장 중요해요(Python 버전을 고정하는 것과 같은 이유).
예를 들어 macOS에서 Homebrew로 JDK를 설치하면 그 java는 $PATH에 없으므로, java_executable을 명시적으로 가리켜요:
[sdk]
coordinators = {
"java-jdk17": {
"classpath": "airflow.sdk.coordinators.java.JavaCoordinator",
"kwargs": {
"jars_root": ["/opt/airflow/jars"],
"java_executable": "/opt/homebrew/opt/openjdk@17/bin/java"
}
}
}
queue_to_coordinator = {"java": "java-jdk17"}
제한 사항
- Task instance당 하나의 JVM 서브프로세스. 각 task instance는 새 JVM을 생성해요. 인스턴스 간 프로세스 내 상태를 공유해야 하는 Task는 대신 XCom이나 외부 저장소를 사용해야 해요.
- assets, deferral, 기타 Airflow 기능에 대한 제한적 지원. 이는 사용자 피드백과 수요에 기반해 향후 구현될 수 있어요.