Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -589,6 +589,46 @@ Durations and date-times are ISO-8601 strings in annotations (``retryDelay = "PT
``java.time.OffsetDateTime`` values in ``config`` calls. An unknown key or a mismatched value type
fails the build for an annotation, and the ``config`` call itself for an object.

.. _java-sdk/native-dag-parsing:

Parsing native Java Dags
~~~~~~~~~~~~~~~~~~~~~~~~

To have Airflow parse the Dags a bundle JAR declares, put the JAR in a Dag bundle and configure a
:class:`~airflow.sdk.coordinators.java.JavaCoordinator` that reads that bundle: leave ``jars_root``
unset to read each Dag's own bundle, or set ``dag_bundle_name``. The Dag processor runs the JAR's main
class to list its Dags, so it needs a Java executable, as the workers do:

.. code-block:: ini

[sdk]
coordinators = {
"java-native": {
"classpath": "airflow.sdk.coordinators.java.JavaCoordinator",
"kwargs": {"java_executable": "/usr/lib/jvm/java-17-openjdk/bin/java"}
}
}
queue_to_coordinator = {"java-native": "java-native"}

A coordinator with ``jars_root`` set only runs tasks; it parses no Dags.

* Every JAR in the bundle whose manifest sets ``Main-Class`` is parsed. Each Dag its main class
declares, through ``Bundle.register`` of a ``DagDef`` or an ``@Builder.Dag`` class, is stored with
that JAR as its file. Task handlers for a Python Dag are not Dags, and a JAR without ``Main-Class``,
such as a dependency of a thin bundle, is skipped.
* Do not declare a Dag in Java that a Python file in the same bundle also defines.
* Keep one executable JAR per bundle, or set ``main_class``, so that only JARs with that
``Main-Class`` are parsed. List JARs that should not be parsed in ``.airflowignore``.
* Set ``queue`` on every task, with ``@Builder.Task(queue = "java-native")`` or
``TaskDef.config("queue", "java-native")``, so it runs on the coordinator's queue. There is no
Dag-level queue yet.
* The Code view shows the source file the Gradle plugin packs into the JAR: the main class by
default, or the file set with ``airflowBundle { dagSource = file("...") }``.
* Cluster policies (``dag_policy``, ``task_policy``) are not applied to a native Java Dag.
* ``airflow dags reserialize`` does not store the Dags of a JAR, which only the Dag processor stores.
``airflow dags test``, ``tasks test``, ``tasks render`` and ``tasks list`` refuse a native Java
Dag.

.. _java-sdk/logging:

Logging
Expand Down Expand Up @@ -855,6 +895,17 @@ The ``build/bundle/`` directory contains all required JAR(s). Copy or mount it i
by ``jars_root`` in the coordinator configuration. :class:`~airflow.sdk.coordinators.java.JavaCoordinator`
scans ``jars_root`` recursively and builds the classpath automatically.

The plugin also packs the source file of ``mainClass`` into the bundle JAR, so the Airflow UI can show
the source of a native Java Dag (see :ref:`java-sdk/native-dag-parsing`). Set ``dagSource`` in
``airflowBundle`` to pack another file instead:

.. code-block:: groovy

airflowBundle {
mainClass = "com.example.Main"
dagSource = file("src/main/java/com/example/MyDag.java")
}

.. note::

You only need the ``annotationProcessor`` entry if you use the annotation-based API. It is not needed for
Expand Down Expand Up @@ -1029,6 +1080,28 @@ this directory.
Unlike the Gradle plugin, Maven has no equivalent of the ``verifyBundleMainClass`` validation step.
A wrong ``<mainClass>`` value will not be caught until runtime.

To show the source of a native Java Dag in the Airflow UI, pack the source file into the bundle JAR and
name its entry with the ``Airflow-Java-SDK-Dag-Code`` manifest attribute:

.. code-block:: xml

<build>
<resources>
<!-- Declaring resources replaces the default, so keep it. -->
<resource>
<directory>src/main/resources</directory>
</resource>
<resource>
<directory>src/main/java/com/example</directory>
<includes><include>Main.java</include></includes>
<targetPath>META-INF/airflow/dag-code/com/example</targetPath>
</resource>
</resources>
</build>

Then add ``<Airflow-Java-SDK-Dag-Code>META-INF/airflow/dag-code/com/example/Main.java</Airflow-Java-SDK-Dag-Code>``
to the ``manifestEntries`` of ``maven-shade-plugin`` or ``maven-jar-plugin`` shown above.

.. _java-sdk/coordinator-config:

:class:`~airflow.sdk.coordinators.java.JavaCoordinator` configuration
Expand Down Expand Up @@ -1065,7 +1138,9 @@ All ``kwargs`` in the ``coordinators`` config entry are passed to the
- Explicit entry-point class. If omitted, the coordinator scans for a JAR whose
manifest sets ``Main-Class`` — in ``jars_root`` when set, otherwise across the
resolved Dag bundle. If multiple executable JARs match the result is
non-deterministic; set ``main_class`` explicitly in that case.
non-deterministic; set ``main_class`` explicitly in that case. When the
coordinator parses native Java Dags, only JARs with this ``Main-Class`` are
parsed.
* - ``task_startup_timeout``
- ``10.0``
- Seconds to wait for the JVM subprocess to connect after launch. Increase this if your
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -22,11 +22,15 @@ package org.apache.airflow.sdk.plugin
import com.github.jengelman.gradle.plugins.shadow.tasks.ShadowJar
import org.gradle.api.Plugin
import org.gradle.api.Project
import org.gradle.api.file.RegularFileProperty
import org.gradle.api.provider.Property
import org.gradle.api.tasks.Copy
import org.gradle.api.tasks.Input
import org.gradle.api.tasks.InputFile
import org.gradle.api.tasks.Optional
import org.gradle.api.tasks.SourceSetContainer
import org.gradle.api.tasks.bundling.Jar
import java.io.File
import java.lang.reflect.Modifier
import java.net.URLClassLoader
import java.util.jar.JarFile
Expand Down Expand Up @@ -61,8 +65,22 @@ abstract class AirflowBundleExtension {
*/
@get:Input
abstract val fatJar: Property<Boolean>

/**
* The Dag source file the Airflow UI shows for this bundle.
*
* It defaults to the `.java` file of [mainClass] in the `main` source set. The file is packed
* into the bundle JAR under `META-INF/airflow/dag-code/`, and the `Airflow-Java-SDK-Dag-Code`
* manifest attribute names it.
*/
@get:InputFile
@get:Optional
abstract val dagSource: RegularFileProperty
}

private const val DAG_CODE_ATTRIBUTE = "Airflow-Java-SDK-Dag-Code"
private const val DAG_CODE_DIR = "META-INF/airflow/dag-code"

/**
* Gradle plugin for building Apache Airflow Java SDK bundles.
*
Expand Down Expand Up @@ -98,6 +116,9 @@ abstract class AirflowBundleExtension {
* directory instead. In this mode, `Airflow-Supervisor-Schema-Version` lives in
* the `airflow-sdk` JAR instead. The bundle JAR still contains `Main-Class`.
*
* In both modes the bundle JAR also carries the Dag source (see
* [AirflowBundleExtension.dagSource]), so the Airflow UI can show the source of a
* native Java Dag.
*/
class AirflowSdkPlugin : Plugin<Project> {
override fun apply(project: Project) {
Expand All @@ -107,12 +128,17 @@ class AirflowSdkPlugin : Plugin<Project> {
ext.fatJar.convention(true)

project.afterEvaluate {
val dagCode = resolveDagCode(project, ext)
project.tasks.withType(Jar::class.java).configureEach { task ->
task.doFirst {
ext.mainClass.orNull?.let { className ->
task.manifest.attributes(mapOf("Main-Class" to className))
}
}
dagCode?.let { (file, entryDir) ->
task.from(file) { spec -> spec.into(entryDir) }
task.manifest.attributes(mapOf(DAG_CODE_ATTRIBUTE to "$entryDir/${file.name}"))
}
}

val classFiles =
Expand Down Expand Up @@ -219,3 +245,35 @@ class AirflowSdkPlugin : Plugin<Project> {
}
}
}

/**
* Returns the Dag source file to pack and its directory in the JAR, or `null` when there is none.
*
* The directory mirrors the package of `mainClass`, so the entry reads like a source path.
*/
private fun resolveDagCode(
project: Project,
ext: AirflowBundleExtension,
): Pair<File, String>? {
val mainClass = ext.mainClass.orNull ?: return null
val className = mainClass.substringBefore('$')
val packagePath = className.substringBeforeLast('.', "").replace('.', '/')
val entryDir = if (packagePath.isEmpty()) DAG_CODE_DIR else "$DAG_CODE_DIR/$packagePath"
ext.dagSource.orNull?.let { return it.asFile to entryDir }

val relativePath = className.replace('.', '/') + ".java"
val source =
project.extensions
.getByType(SourceSetContainer::class.java)
.getByName("main")
.java.srcDirs
.map { it.resolve(relativePath) }
.firstOrNull { it.isFile }
if (source == null) {
project.logger.info(
"airflowBundle: no $relativePath in the main source set, so the bundle JAR carries no Dag " +
"source for the Airflow UI. Set airflowBundle.dagSource to choose the file.",
)
}
return source?.let { it to entryDir }
}
86 changes: 86 additions & 0 deletions task-sdk/src/airflow/sdk/coordinators/java/_dag_importer.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,86 @@
#
# Licensed to the Apache Software Foundation (ASF) under one
# or more contributor license agreements. See the NOTICE file
# distributed with this work for additional information
# regarding copyright ownership. The ASF licenses this file
# to you under the Apache License, Version 2.0 (the
# "License"); you may not use this file except in compliance
# with the License. You may obtain a copy of the License at
#
# http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing,
# software distributed under the License is distributed on an
# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
# KIND, either express or implied. See the License for the
# specific language governing permissions and limitations
# under the License.
"""The Dag importer of :class:`~airflow.sdk.coordinators.java.JavaCoordinator`."""

from __future__ import annotations

import zipfile
from typing import TYPE_CHECKING, ClassVar, Final

from airflow.sdk.coordinators._dag_importer import CoordinatorDagImporter
from airflow.sdk.coordinators.java._jar_manifest import DAG_CODE, MAIN_CLASS, read_main_attributes
from airflow.sdk.importers.base import DagSourceCode

if TYPE_CHECKING:
from airflow.sdk.coordinators.java.coordinator import JavaCoordinator
from airflow.sdk.importers.base import DagDefinition

_MAX_SOURCE_BYTES: Final = 1024 * 1024
_NO_SOURCE: Final = (
"// This JAR embeds no Dag source. Build it with the Airflow Java SDK Gradle plugin, or set\n"
"// airflowBundle.dagSource, to show the source here.\n"
)
_SOURCE_TOO_LARGE: Final = "// The Dag source this JAR embeds is over 1 MiB, so it is not shown.\n"


class JavaDagImporter(CoordinatorDagImporter):
"""
Parse the native Dags of Java bundle JARs with the coordinator's JVM.

Only a JAR whose manifest sets ``Main-Class`` is parsed, so the dependency JARs of a thin bundle
are skipped.
"""

artifact_suffix: ClassVar[str] = ".jar"
supported_extensions = [".jar"]
coordinator: JavaCoordinator

def might_contain_dag(self, definition: DagDefinition, safe_mode: bool) -> bool:
"""
Return whether the JAR sets ``Main-Class``, matching the coordinator's ``main_class`` if that is set.

``safe_mode`` does not apply, because a JAR without ``Main-Class`` cannot run at all. A JAR that
cannot be read is kept, so that parsing it reports the error.
"""
try:
with definition.as_file() as path, zipfile.ZipFile(path) as zf:
attributes = read_main_attributes(zf) or {}
except (OSError, zipfile.BadZipFile):
return True
if not (main_class := attributes.get(MAIN_CLASS)):
return False
return not self.coordinator.main_class or main_class == self.coordinator.main_class

def get_source_code(self, definition: DagDefinition) -> DagSourceCode:
"""Return the Dag source the JAR embeds, or a placeholder when it embeds none."""
with definition.as_file() as path, zipfile.ZipFile(path) as zf:
if (info := _find_source_entry(zf)) is None:
return DagSourceCode(source_code=_NO_SOURCE, language="java")
if info.file_size > _MAX_SOURCE_BYTES:
return DagSourceCode(source_code=_SOURCE_TOO_LARGE, language="java")
source = zf.read(info).decode("utf-8", errors="replace")
return DagSourceCode(source_code=source or _NO_SOURCE, language="java")


def _find_source_entry(zf: zipfile.ZipFile) -> zipfile.ZipInfo | None:
if not (entry := (read_main_attributes(zf) or {}).get(DAG_CODE)):
return None
try:
return zf.getinfo(entry)
except KeyError:
return None
71 changes: 71 additions & 0 deletions task-sdk/src/airflow/sdk/coordinators/java/_jar_manifest.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,71 @@
#
# Licensed to the Apache Software Foundation (ASF) under one
# or more contributor license agreements. See the NOTICE file
# distributed with this work for additional information
# regarding copyright ownership. The ASF licenses this file
# to you under the Apache License, Version 2.0 (the
# "License"); you may not use this file except in compliance
# with the License. You may obtain a copy of the License at
#
# http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing,
# software distributed under the License is distributed on an
# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
# KIND, either express or implied. See the License for the
# specific language governing permissions and limitations
# under the License.
"""Read JAR manifests as the JAR File Specification defines them."""

from __future__ import annotations

import re
from typing import TYPE_CHECKING, Final

if TYPE_CHECKING:
import zipfile

MANIFEST_NAME: Final = "META-INF/MANIFEST.MF"

# Attribute keys as parse_main_attributes returns them, lower-cased.
MAIN_CLASS: Final = "main-class"
SUPERVISOR_SCHEMA_VERSION: Final = "airflow-supervisor-schema-version"
DAG_CODE: Final = "airflow-java-sdk-dag-code"

_LINE_END = re.compile(rb"\r\n|\r|\n")


def parse_main_attributes(data: bytes) -> dict[str, str]:
"""
Return the main-section attributes of a JAR manifest, keyed by lower-cased name.

Attribute names are case-insensitive. Lines end with CRLF, LF or CR, and a line that starts with
one space continues the previous value. The main section ends at the first blank line.
"""
headers: list[bytearray] = []
for line in _LINE_END.split(data):
if not line:
break
if line.startswith(b" "):
# Lines are folded by bytes, which can split a multi-byte character, so unfold
# before decoding.
if headers:
headers[-1] += line[1:]
continue
headers.append(bytearray(line))

attributes: dict[str, str] = {}
for header in headers:
name, sep, value = header.decode("utf-8", errors="replace").partition(":")
if sep:
attributes[name.strip().lower()] = value.removeprefix(" ")
return attributes


def read_main_attributes(zf: zipfile.ZipFile) -> dict[str, str] | None:
"""Return the main attributes of the JAR open in *zf*, or ``None`` when it has no manifest."""
try:
info = zf.getinfo(MANIFEST_NAME)
except KeyError:
return None
return parse_main_attributes(zf.read(info))
Loading
Loading