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:

Reply via email to