Repository navigation
Conversation
…explosion DagModelOperation.find_orm_dags issued a single query with five joinedload() calls on DagModel's one-to-many collections (tags, schedule_asset_references, schedule_asset_alias_references, task_outlet_asset_references, dag_owner_links). Multiple joinedload() calls on one-to-many collections in one query produce a cartesian product: a production deployment observed 500 input dag_ids returning 3,907 rows. The database executed the query quickly (16ms), but the client had to receive, deserialize, and de-duplicate every exploded row, including redundant copies of DagModel's JSON columns (partition_mapper_info, asset_expression, deadline) per duplicate row. This was directly responsible for multi-second client-side delays in production. Switch all five relationships to selectinload, which issues one follow-up "WHERE dag_id IN (...)" query per collection instead of joining them all at once, eliminating the row multiplication. DagModelOperation.find_orm_dags is called twice per DAG.bulk_write_to_db invocation, so this trades 2 statements for up to 12 (1 base + 5 selectin, twice) -- more round trips, but each one returns only real rows, and no client-side dedup of exploded duplicates. Also drop innerjoin=False from the tags load (a joinedload-specific knob to avoid dropping DAGs with no tags via an inner join; selectinload has no such concern since it's a separate query keyed off the base result), and update the pinned per-call statement budget in test_manager.py to account for the extra round trips. Verified locally against a synthetic 400-DAG dynamically-generated file (5 tags + 2 asset outlets + 2 owner links each): the unpatched joinedload query returns 8,000 rows for the 400 dag_ids (a 20x explosion), while the patched version returns exactly the right rows with no explosion. Steady-state re-sync of that file dropped from ~0.75s to ~0.54s locally against Postgres on localhost (near-zero network latency); the gap should be far larger for any real client with actual network latency and deserialization overhead, which is the scenario this issue describes. Closes apache#72393 Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
SameerMesiah97
left a comment
There was a problem hiding this comment.
Given this caused a real-world performance issue, would it also be worth adding a focused regression test around find_orm_dags() that verifies the five relationships are loaded without multi-collection joins?
|
test failure is unrelated: can push a dummy commit to rerun or just let it stand, I'll wait for the rest of CI to finish either way. |
I would rebase or just convert to draft. That re-triigers CI. I would recommend rebase though. |
…oad change CI caught a separate pinned query-count assertion this PR's earlier fix to test_manager.py's FIXED_PER_CALL missed: TestDag::test_bulk_write_to_db, test_bulk_write_to_db_single_dag, and test_bulk_write_to_db_multiple_dags in airflow-core/tests/unit/models/test_dag.py each pin an exact assert_queries_count() for SerializedDAG.bulk_write_to_db(), which calls DagModelOperation.find_orm_dags() twice per invocation. Verified the real counts empirically against Postgres (via a small script replicating each test's assert_queries_count() call sites with count_queries() instead, since the assertion only raises on overage rather than reporting the actual number): - First (insert) call, no pre-existing DagModel rows for the initial find_orm_dags(), but the "refetch" call after flushing new assets does match rows and pays the 5 selectinload follow-ups: 6 -> 10 in all three tests. - Steady-state re-sync, where both find_orm_dags() calls match existing DagModels and both pay the full 5 selectinload follow-ups: 9 -> 14 (test_bulk_write_to_db, test_bulk_write_to_db_multiple_dags) and 8 -> 14 (test_bulk_write_to_db_single_dag). - Adding/removing tags in test_bulk_write_to_db, same steady-state shape: 10 -> 15. Ran the full test_dag.py (209 tests) plus test_collection.py/test_manager.py again against Postgres -- all green. Follow-up to GH#72393 / the collection.py selectinload change. Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
tests were indeed broken (just not that one) corrected and passing now |
Address Sameer's review on apache#72395: - Add TestFindOrmDagsLoaderStrategy to test_collection.py with two tests that guard the actual selectinload behavior, not just a query count that could be "corrected" back down without anyone noticing the row-explosion came back: - test_find_orm_dags_uses_selectinload inspects the loader options attached to find_orm_dags()'s compiled statement directly, and fails if any of the five relationships (tags, schedule_asset_references, schedule_asset_alias_references, task_outlet_asset_references, dag_owner_links) uses anything other than "selectin". - test_find_orm_dags_does_not_multiply_rows creates one Dag with 3 tags, 2 owner links, and 2 task outlet asset references, then executes find_orm_dags()'s captured statement at the raw DBAPI cursor level (bypassing the ORM's own .unique() row processing, which would otherwise hide exactly the multiplication this test exists to catch) and asserts the base query returns exactly 1 row. Verified both tests actually catch the regression: manually reverted find_orm_dags to the old multi-collection joinedload() and confirmed both fail (the second one reporting the expected 3x2x2=12 raw rows for that one dag_id), then restored the selectinload fix and confirmed both pass again. - Tightened the code comments on collection.py's find_orm_dags() and test_manager.py's FIXED_PER_CALL per review, and did the same for the related comments added to test_dag.py in the prior commit. Ran the full test_dag.py (209), test_collection.py (66 passed, 2 unrelated FAB skips), and test_manager.py (165) suites again against Postgres -- all green. Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
|
@SameerMesiah97 addressed both — thanks for the review. Regression test added (
Verified both actually catch the regression: temporarily reverted to the old multi- Comments tightened on All green on Postgres: |
|
Hello @seanmuth - thank you for your contributions to Apache Airflow! The Airflow community has introduced a limit of 5 open pull requests at a time for contributors without write access to the repository. You currently have 7 open pull requests, so - as a one-time step of introducing the limit - we closed the ones where maintainers have not engaged yet:
These pull requests stay open because maintainers are already engaged in them - they count towards your limit:
This is not a judgement of you or of your changes. We never told contributors before that opening many pull requests at once was a problem, so there is nothing to feel bad about - and nothing is lost: your branches, commits and the review history stay where they are. What we ask you to do is to make your first prioritization decision: choose which of the pull requests above matter most to you, and reopen them (up to 5 open at a time, including the ones still open) with the "Reopen pull request" button or While your pull requests are waiting for review, the most valuable thing you can do is help in other ways - reviewing other contributors' pull requests, helping with issues, and taking part in the discussions on the devlist and Slack. Why we introduced the limit, what it means for you and how to reopen or restore a pull request is explained in https://github.com/apache/airflow/blob/main/contributing-docs/32_open_pull_request_limit.rst. Drafted-by: Claude Code (Opus 5); reviewed by @potiuk before posting |
closes: #72393
Problem
DagModelOperation.find_orm_dagsinairflow-core/src/airflow/dag_processing/collection.pyeagerly loads five different one-to-many collections onDagModel(tags,schedule_asset_references,schedule_asset_alias_references,task_outlet_asset_references,dag_owner_links) usingjoinedload()in a single query:Combining multiple
joinedload()calls on one-to-many collections into a single query produces a cartesian product across those collections. This was confirmed live on a production deployment: 500 inputdag_ids produced 3,907 result rows viaEXPLAIN (ANALYZE, BUFFERS). The database itself executed the query quickly (16ms) — the real cost lands on the client, which has to receive, deserialize, and de-duplicate (via.unique()) every exploded row, including redundant copies ofDagModel's JSON columns (partition_mapper_info,asset_expression,deadline) once per duplicate. This was observed causing multi-second processing delays per call in production for an async SQLAlchemy/asyncpg client, but it affects any client — sync or async — syncing a DAG set with a meaningful number of tags, asset references, or owner links.Fix
Switch all five relationships to
selectinload(). Instead of one big join, this issues one follow-upWHERE dag_id IN (...)query per collection — the same eager-loading outcome (all five collections populated on the returnedDagModelobjects), but each query returns only rows that actually exist, with no multiplication.This is a real tradeoff worth being explicit about:
find_orm_dagsis called twice perDAG.bulk_write_to_db()invocation (once to look up existing DagModels, once to refetch after flushing new assets), so this goes from 2 statements total (1 per call) to up to 12 (1 base + 5 selectin, times 2 calls). More round trips, but each one is cheap, returns exactly the right number of rows, and avoids the client-side deserialize/dedup cost of the exploded joined result.Also dropped
innerjoin=Falsefrom thetagsload — that's ajoinedload-specific knob to prevent an inner join from silently dropping DAGs with no tags.selectinloadhas no equivalent concern since it's a separate query keyed off the dag_ids already returned by the base select, so a DAG with zero tags is unaffected either way.Testing
airflow-core/tests/unit/dag_processing/test_collection.pyandtest_manager.pysuites against both sqlite and Postgres backends — all passing.test_manager.pyhas a pinned per-call SQL statement budget (FIXED_PER_CALL) that tracks exactly this kind of change. Updated it from 9 to 19 with an inline comment explaining the +10 (5 extraselectinstatements perfind_orm_dags()call, times 2 calls per persistence call), so a future statement-count regression there is easy to diagnose rather than a mystery.globals()[dag_id] = dagpattern), each with 5 tags, 2 task-level asset outlets, and 2 owner links, then ran it throughSerializedDAG.bulk_write_to_db(the same path the scheduler's dag processor uses) against a local Postgres instance.joinedloadquery returns 8,000 rows for those 400dag_ids (a 20x explosion from the 5×2×2 collection sizes) — directionally consistent with, and worse than, the 500→3,907 (~7.8x) production case.🤖 Generated with Claude Code