liwenjie200543 commented on code in PR #73911:
URL: https://github.com/apache/airflow/pull/73911#discussion_r4140180491


##########
airflow-core/tests/unit/dag_processing/test_collection.py:
##########
@@ -1565,6 +1565,71 @@ def test_new_errors_keep_their_own_bundle_name(self, 
_mock_full_path, session):
         assert rows == {("bundle_a", "a.py"), ("bundle_b", "b.py")}
 
 
[email protected]_test
+class TestFindOrmDagsEagerLoading:
+    """find_orm_dags must not eager-join one-to-many collections in the main 
query.
+
+    Joining several one-to-many collections with joinedload in one statement
+    multiplies the returned rows (3 tags x 2 owner links => 6 copies of the
+    same wide dag row), which dominated DagModel sync time on large
+    deployments (#72393). The loader strategy is asserted via the SQL that
+    actually reaches the database: no statement may join the one-to-many
+    tables; they are loaded through selectinload's secondary selects instead.
+    """
+
+    DAG_ID = "test_dag"
+
+    @pytest.fixture(autouse=True)
+    def setup_teardown(self, session):
+        yield
+        session.execute(delete(DagModel).where(DagModel.dag_id == self.DAG_ID))
+        session.commit()
+
+    @pytest.fixture
+    def capture_statements(self, session):
+        statements = []
+
+        def before_cursor_execute(conn, cursor, statement, parameters, 
context, executemany):
+            statements.append(statement)
+
+        event.listen(session.bind, "before_cursor_execute", 
before_cursor_execute)
+        yield statements
+        event.remove(session.bind, "before_cursor_execute", 
before_cursor_execute)
+
+    def test_find_orm_dags_does_not_eager_join_one_to_many_collections(
+        self, dag_maker, session, capture_statements
+    ):
+        from airflow.models.dag import DagOwnerAttributes
+
+        with dag_maker(self.DAG_ID, schedule=None):
+            EmptyOperator(task_id="empty")
+        dag_maker.sync_dagbag_to_db()
+
+        dag_model = session.scalars(select(DagModel).where(DagModel.dag_id == 
self.DAG_ID)).one()
+        dag_model.tags = [DagTag(name=f"tag-{i}", dag_id=self.DAG_ID) for i in 
range(3)]
+        session.add_all(
+            DagOwnerAttributes(dag_id=self.DAG_ID, owner=f"owner-{i}", 
link="https://example.com";)
+            for i in range(2)
+        )
+        session.commit()
+        session.expire_all()
+
+        dags = {self.DAG_ID: LazyDeserializedDAG.from_dag(dag_maker.dag)}
+        orm_dags = DagModelOperation(dags, "testing", 
None).find_orm_dags(session=session)
+
+        assert set(orm_dags) == {self.DAG_ID}
+        assert {t.name for t in orm_dags[self.DAG_ID].tags} == {"tag-0", 
"tag-1", "tag-2"}
+        assert {o.owner for o in orm_dags[self.DAG_ID].dag_owner_links} == 
{"owner-0", "owner-1"}
+        joined_one_to_many = [
+            s
+            for s in capture_statements
+            if any(join in s for join in ("JOIN dag_tag", "JOIN 
dag_owner_attributes"))
+        ]
+        assert joined_one_to_many == [], (
+            "one-to-many collections must be loaded via selectinload, not 
joinedload"
+        )

Review Comment:
   Thanks for the review! I've strengthened the tests:
   
   - The captured statements are now cleared right before calling 
`find_orm_dags()`, and the tests assert that accessing each of the five 
collections afterwards executes **no additional SQL** — so lazy loading can no 
longer pass (verified: with the `selectinload` options removed, both tests fail 
with the "executed additional SQL" assertion).
   - Added a test covering a DAG with no tags, owner links, or asset 
references, asserting it is returned with all five collections empty and no 
further SQL.
   - Extended the content assertions to `schedule_asset_references`, 
`schedule_asset_alias_references`, and `task_outlet_asset_references` (asserted 
via their FK columns — touching `.asset` would itself lazy-load and is outside 
`find_orm_dags`' loading scope).



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to