Skip to content

Use selectinload instead of joinedload in find_orm_dags to avoid row explosion - #72395

Closed
seanmuth wants to merge 3 commits into
apache:mainfrom
seanmuth:fix-find-orm-dags-row-explosion-72393
Closed

seanmuth wants to merge 3 commits into
apache:mainfrom
seanmuth:fix-find-orm-dags-row-explosion-72393

Conversation

@seanmuth

@seanmuth seanmuth commented Sep 1, 2026

Copy link
Copy Markdown
Contributor

closes: #72393

Problem

DagModelOperation.find_orm_dags in airflow-core/src/airflow/dag_processing/collection.py eagerly loads five different one-to-many collections on DagModel (tags, schedule_asset_references, schedule_asset_alias_references, task_outlet_asset_references, dag_owner_links) using joinedload() in a single query:

stmt = with_row_locks(
    (
        select(DagModel)
        .options(joinedload(DagModel.tags, innerjoin=False))
        .where(DagModel.dag_id.in_(self.dags))
        .options(joinedload(DagModel.schedule_asset_references))
        .options(joinedload(DagModel.schedule_asset_alias_references))
        .options(joinedload(DagModel.task_outlet_asset_references))
        .options(joinedload(DagModel.dag_owner_links))
    ),
    of=DagModel,
    session=session,
)

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 input dag_ids produced 3,907 result rows via EXPLAIN (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 of DagModel'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-up WHERE dag_id IN (...) query per collection — the same eager-loading outcome (all five collections populated on the returned DagModel objects), but each query returns only rows that actually exist, with no multiplication.

This is a real tradeoff worth being explicit about: find_orm_dags is called twice per DAG.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=False from the tags load — that's a joinedload-specific knob to prevent an inner join from silently dropping DAGs with no tags. selectinload has 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

  • Ran the full airflow-core/tests/unit/dag_processing/test_collection.py and test_manager.py suites against both sqlite and Postgres backends — all passing.
  • test_manager.py has 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 extra selectin statements per find_orm_dags() call, times 2 calls per persistence call), so a future statement-count regression there is easy to diagnose rather than a mystery.
  • Local verification: generated a synthetic file with 400 dynamically-generated DAGs (globals()[dag_id] = dag pattern), each with 5 tags, 2 task-level asset outlets, and 2 owner links, then ran it through SerializedDAG.bulk_write_to_db (the same path the scheduler's dag processor uses) against a local Postgres instance.
    • Confirmed the unpatched joinedload query returns 8,000 rows for those 400 dag_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.
    • Confirmed the patched version persists all 400 DAGs correctly — tags, owner links, and task outlet asset references all present and correct on spot-checked DAGs — with no explosion.
    • Timed steady-state re-sync of the same 400-DAG file 5x each way on the same local Postgres instance (near-zero network latency, so this understates the real-world win): unpatched avg ~0.75s, patched avg ~0.54s (~28% faster). The gap should be substantially larger for any real deployment with actual client-DB network latency and per-row deserialization overhead, which is exactly the scenario described in the issue.

🤖 Generated with Claude Code

…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 SameerMesiah97 left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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?

Comment thread airflow-core/src/airflow/dag_processing/collection.py Outdated
Comment thread airflow-core/tests/unit/dag_processing/test_manager.py Outdated
@seanmuth

seanmuth commented Sep 1, 2026

Copy link
Copy Markdown
Contributor Author

test failure is unrelated:

Run prek install-hooks
  error: Failed to install hook `go-mod-tidy`
    caused by: Failed to install go
    caused by: Failed to download go
    caused by: Failed to download file from https://go.dev/dl/go1.25.0.linux-amd64.tar.gz
    caused by: error sending request for url (https://go.dev/dl/go1.25.0.linux-amd64.tar.gz)
    caused by: client error (Connect)
    caused by: Connection reset by peer (os error 104)
  Error: Process completed with exit code 2.

can push a dummy commit to rerun or just let it stand, I'll wait for the rest of CI to finish either way.

@SameerMesiah97

SameerMesiah97 commented Sep 1, 2026 •

Copy link
Copy Markdown
Contributor

test failure is unrelated:

Run prek install-hooks
  error: Failed to install hook `go-mod-tidy`
    caused by: Failed to install go
    caused by: Failed to download go
    caused by: Failed to download file from https://go.dev/dl/go1.25.0.linux-amd64.tar.gz
    caused by: error sending request for url (https://go.dev/dl/go1.25.0.linux-amd64.tar.gz)
    caused by: client error (Connect)
    caused by: Connection reset by peer (os error 104)
  Error: Process completed with exit code 2.

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>
@seanmuth

seanmuth commented Sep 2, 2026

Copy link
Copy Markdown
Contributor Author

test failure is unrelated:

Run prek install-hooks
  error: Failed to install hook `go-mod-tidy`
    caused by: Failed to install go
    caused by: Failed to download go
    caused by: Failed to download file from https://go.dev/dl/go1.25.0.linux-amd64.tar.gz
    caused by: error sending request for url (https://go.dev/dl/go1.25.0.linux-amd64.tar.gz)
    caused by: client error (Connect)
    caused by: Connection reset by peer (os error 104)
  Error: Process completed with exit code 2.

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.

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>
@seanmuth

seanmuth commented Sep 2, 2026

Copy link
Copy Markdown
Contributor Author

@SameerMesiah97 addressed both — thanks for the review.

Regression test added (TestFindOrmDagsLoaderStrategy in test_collection.py): two tests that check the actual mechanism, not just a statement count that could be blindly "corrected" back:

  • test_find_orm_dags_uses_selectinload — inspects the compiled statement's ORM loader options directly, fails if any of the five relationships isn't selectin.
  • test_find_orm_dags_does_not_multiply_rows — a DAG with 3 tags × 2 owner links × 2 asset refs, checked at the raw DBAPI cursor level (bypassing .unique(), which would otherwise hide the exact multiplication this needs to catch) — asserts exactly 1 raw row for 1 dag_id.

Verified both actually catch the regression: temporarily reverted to the old multi-joinedload() code and confirmed both tests fail correctly ('joined' instead of 'selectin', and 12 raw rows — 3×2×2 — instead of 1), then restored the fix and confirmed green.

Comments tightened on collection.py and test_manager.py per your suggested wording.

All green on Postgres: test_collection.py (66 passed), test_manager.py (165 passed), test_dag.py (209 passed).

@potiuk potiuk added the closed because of open PR limit Closed as a one-time step of introducing the open pull request limit label Sep 25, 2026
@potiuk

potiuk commented Sep 25, 2026

Copy link
Copy Markdown
Member

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 gh pr reopen <PR_NUMBER> --repo apache/airflow. Reopen the ones you are ready to follow through - keep them rebased, respond to review comments and fix failing checks.

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

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:dag-processor closed because of open PR limit Closed as a one-time step of introducing the open pull request limit

Projects

None yet

Development

Successfully merging this pull request may close these issues.

find_orm_dags: multiple joinedload() collections cause row-explosion, slow DagModel sync (commonly >10s)

3 participants