This is an automated email from the ASF dual-hosted git repository.
jason810496 pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/airflow.git
The following commit(s) were added to refs/heads/main by this push:
new 1219b1a8541 Eager-load DagCode in the Dag source API (#73752)
1219b1a8541 is described below
commit 1219b1a8541e06b46450c8c99e0a3bef92a2be91
Author: Tzu-ping Chung <[email protected]>
AuthorDate: Sat Sep 26 23:37:43 2026 +0800
Eager-load DagCode in the Dag source API (#73752)
GET /dags/{dag_id}/dagSources loaded the DagVersion first and then read
the lazy dag_code relationship, issuing a second query. Add a
load_dag_code option to DagVersion.get_version that joinedloads dag_code,
and use it from the Dag source route so the source and its language come
back in a single query.
Closes: #73716
Signed-off-by: Tzu-ping Chung <[email protected]>
---
.../core_api/routes/public/dag_sources.py | 2 +-
airflow-core/src/airflow/models/dag_version.py | 5 +++++
airflow-core/tests/unit/models/test_dag_version.py | 24 ++++++++++++++++++++++
3 files changed, 30 insertions(+), 1 deletion(-)
diff --git
a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/dag_sources.py
b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/dag_sources.py
index 6723938ec1b..f74ce76c7b7 100644
--- a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/dag_sources.py
+++ b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/dag_sources.py
@@ -63,7 +63,7 @@ def get_dag_source(
version_number: int | None = None,
):
"""Get source code using file token."""
- dag_version = DagVersion.get_version(dag_id, version_number,
session=session)
+ dag_version = DagVersion.get_version(dag_id, version_number,
load_dag_code=True, session=session)
if not dag_version:
raise HTTPException(
status.HTTP_404_NOT_FOUND,
diff --git a/airflow-core/src/airflow/models/dag_version.py
b/airflow-core/src/airflow/models/dag_version.py
index c2abb6e3a26..ff4971f4ba8 100644
--- a/airflow-core/src/airflow/models/dag_version.py
+++ b/airflow-core/src/airflow/models/dag_version.py
@@ -233,6 +233,7 @@ class DagVersion(Base):
dag_id: str,
version_number: int | None = None,
*,
+ load_dag_code: bool = False,
session: Session = NEW_SESSION,
) -> DagVersion | None:
"""
@@ -241,6 +242,7 @@ class DagVersion(Base):
:param dag_id: The DAG ID.
:param version_number: The version number to look up. When ``None``,
the latest
version is returned; any other value -- ``0`` included -- is used
as a filter.
+ :param load_dag_code: Whether to eagerly load the associated DagCode
in the same query.
:param session: The database session.
:return: The version of the DAG or None if not found.
"""
@@ -248,6 +250,9 @@ class DagVersion(Base):
if version_number is not None:
version_select_obj = version_select_obj.where(cls.version_number
== version_number)
+ if load_dag_code:
+ version_select_obj =
version_select_obj.options(joinedload(cls.dag_code))
+
return
session.scalar(version_select_obj.order_by(cls.version_number.desc()).limit(1))
@property
diff --git a/airflow-core/tests/unit/models/test_dag_version.py
b/airflow-core/tests/unit/models/test_dag_version.py
index 176378691d3..fd187a4935c 100644
--- a/airflow-core/tests/unit/models/test_dag_version.py
+++ b/airflow-core/tests/unit/models/test_dag_version.py
@@ -166,6 +166,30 @@ class TestDagVersion:
assert version.dag_id == dag1_id
assert version.version == f"{dag1_id}-1"
+ @pytest.mark.need_serialized_dag
+ def test_get_version_eager_loads_dag_code(self, dag_maker, session):
+ """``load_dag_code`` fetches DagCode in the same query, so accessing
it costs no round trip."""
+ dag_id = "test_eager_dag_code"
+ with dag_maker(dag_id):
+ EmptyOperator(task_id="task1")
+
+ session.expunge_all()
+ version = DagVersion.get_version(dag_id, load_dag_code=True,
session=session)
+ with assert_queries_count(0):
+ assert version.dag_code.source_code is not None
+
+ @pytest.mark.need_serialized_dag
+ def test_get_version_lazy_loads_dag_code_by_default(self, dag_maker,
session):
+ """Without ``load_dag_code`` the relationship stays lazy and is
fetched on first access."""
+ dag_id = "test_lazy_dag_code"
+ with dag_maker(dag_id):
+ EmptyOperator(task_id="task1")
+
+ session.expunge_all()
+ version = DagVersion.get_version(dag_id, session=session)
+ with assert_queries_count(1):
+ assert version.dag_code.source_code is not None
+
@pytest.mark.need_serialized_dag
def test_version_property(self, dag_maker):
with dag_maker("test1") as dag: