This is an automated email from the ASF dual-hosted git repository. jason810496 pushed a commit to branch jason/lang-sdk-e2e/08-native-e2e in repository https://gitbox.apache.org/repos/asf/airflow.git
commit 9582bdaaf284e1528a699f599965e13a4ac1921a Author: ZHE YOU LIU <[email protected]> AuthorDate: Mon Sep 28 09:02:43 2026 +0000 Share the Dag run lookup and the native Java queue name wait_for_dag_run now polls through get_dag_run, and the annotation Dag takes its queue from the interface Dag's constant. --- .../src/java/org/apache/airflow/e2e/NativeBundleBuilder.java | 2 +- .../java/org/apache/airflow/e2e/nativedag/AnnotationDag.java | 12 +++++++----- .../tests/airflow_e2e_tests/e2e_test_utils/clients.py | 7 ++----- .../java_sdk_tests/test_java_sdk_native_dag.py | 7 ++++++- 4 files changed, 16 insertions(+), 12 deletions(-) diff --git a/airflow-e2e-tests/java-native-bundle/src/java/org/apache/airflow/e2e/NativeBundleBuilder.java b/airflow-e2e-tests/java-native-bundle/src/java/org/apache/airflow/e2e/NativeBundleBuilder.java index 17ac91e4cd3..cfddd691de2 100644 --- a/airflow-e2e-tests/java-native-bundle/src/java/org/apache/airflow/e2e/NativeBundleBuilder.java +++ b/airflow-e2e-tests/java-native-bundle/src/java/org/apache/airflow/e2e/NativeBundleBuilder.java @@ -31,7 +31,7 @@ import org.apache.airflow.sdk.*; * {@code queue}, which routes it to the {@code java-native} coordinator. */ public class NativeBundleBuilder { - private static final String QUEUE = "java-native"; + public static final String QUEUE = "java-native"; public static class Extract implements Task { @Override diff --git a/airflow-e2e-tests/java-native-bundle/src/java/org/apache/airflow/e2e/nativedag/AnnotationDag.java b/airflow-e2e-tests/java-native-bundle/src/java/org/apache/airflow/e2e/nativedag/AnnotationDag.java index e00dccb27b5..89abab7737b 100644 --- a/airflow-e2e-tests/java-native-bundle/src/java/org/apache/airflow/e2e/nativedag/AnnotationDag.java +++ b/airflow-e2e-tests/java-native-bundle/src/java/org/apache/airflow/e2e/nativedag/AnnotationDag.java @@ -20,34 +20,36 @@ // "native" is a Java keyword, so the native Dags live in "nativedag". package org.apache.airflow.e2e.nativedag; +import static org.apache.airflow.e2e.NativeBundleBuilder.QUEUE; + import org.apache.airflow.sdk.*; -/** A native Java Dag declared with annotations; its tasks set {@code queue} like the interface Dag's. */ +/** A native Java Dag declared with annotations; its tasks use the interface Dag's queue. */ @Builder.Dag( id = "java_native_annotation_e2e", description = "Native Java Dag of the Airflow E2E tests, declared with annotations", catchup = false, tags = {"java-sdk", "native", "e2e"}) public class AnnotationDag { - @Builder.Task(id = "extract", queue = "java-native") + @Builder.Task(id = "extract", queue = QUEUE) public long extract() { return 42L; } - @Builder.Task(id = "transform", queue = "java-native") + @Builder.Task(id = "transform", queue = QUEUE) public long transform(long extracted, double factor) { return (long) (extracted * factor); } /** Fails the run unless the value reached it, so a successful run proves the XCom flow. */ - @Builder.Task(id = "load", queue = "java-native") + @Builder.Task(id = "load", queue = QUEUE) public void load(long transformed) { if (transformed != 63L) { throw new IllegalStateException("load expected 63 from transform, got " + transformed); } } - @Builder.Task(id = "audit", queue = "java-native") + @Builder.Task(id = "audit", queue = QUEUE) public void audit() {} @Builder.Deps diff --git a/airflow-e2e-tests/tests/airflow_e2e_tests/e2e_test_utils/clients.py b/airflow-e2e-tests/tests/airflow_e2e_tests/e2e_test_utils/clients.py index ee541f97f7d..bd80556b586 100644 --- a/airflow-e2e-tests/tests/airflow_e2e_tests/e2e_test_utils/clients.py +++ b/airflow-e2e-tests/tests/airflow_e2e_tests/e2e_test_utils/clients.py @@ -144,11 +144,7 @@ class AirflowClient: def wait_for_dag_run(self, dag_id: str, run_id: str, timeout=300, check_interval=5): start_time = time.time() while time.time() - start_time < timeout: - response = self._make_request( - method="GET", - endpoint=f"dags/{dag_id}/dagRuns/{run_id}", - ) - state = response.get("state") + state = self.get_dag_run(dag_id=dag_id, run_id=run_id).get("state") if state in {"success", "failed"}: return state time.sleep(check_interval) @@ -184,6 +180,7 @@ class AirflowClient: return self._make_request(method="GET", endpoint=f"dagSources/{dag_id}") def get_dag_run(self, dag_id: str, run_id: str): + """Get a Dag run, with its state, run type and conf.""" return self._make_request(method="GET", endpoint=f"dags/{dag_id}/dagRuns/{run_id}") def trigger_dag_and_wait(self, dag_id: str, json=None): diff --git a/airflow-e2e-tests/tests/airflow_e2e_tests/java_sdk_tests/test_java_sdk_native_dag.py b/airflow-e2e-tests/tests/airflow_e2e_tests/java_sdk_tests/test_java_sdk_native_dag.py index f8c00b31fb3..f7296c1f96b 100644 --- a/airflow-e2e-tests/tests/airflow_e2e_tests/java_sdk_tests/test_java_sdk_native_dag.py +++ b/airflow-e2e-tests/tests/airflow_e2e_tests/java_sdk_tests/test_java_sdk_native_dag.py @@ -123,7 +123,12 @@ def test_every_task_is_routed_to_the_native_coordinator(parsed_dags: AirflowClie @_by_dag_id def test_the_dag_source_is_the_bundle_main_class(parsed_dags: AirflowClient, native_dag: _NativeDag): - """The Code view shows the Java source the JAR embeds, not the JAR read as text.""" + """ + The Code view shows the Java source the JAR embeds, not the JAR read as text. + + A JAR embeds one source, its main class, so both Dags show the file that declares + ``java_native_e2e`` and registers ``java_native_annotation_e2e``. + """ content = parsed_dags.get_dag_source(native_dag.dag_id)["content"] assert "public class NativeBundleBuilder" in content
