Skip to content

Fix memory issue: remove eager loading of all TIs in scheduler - #60956

Merged
kaxil merged 1 commit into
apache:mainfrom
astronomer:fix-scheduler-ti-loading-memory
Jan 22, 2026
Merged

kaxil merged 1 commit into
apache:mainfrom
astronomer:fix-scheduler-ti-loading-memory

Conversation

@kaxil

@kaxil kaxil commented Jan 22, 2026 •

Copy link
Copy Markdown
Member

The scheduler's get_running_dag_runs_to_examine() was using joinedload(DagRun.task_instances) which loaded ALL TaskInstances for each DagRun into memory. This happens every scheduler loop and could cause significant memory pressure with large DAGs.

Changes:

  • Remove unnecessary joinedload(task_instances) from the query
  • Convert the TI loop in _verify_integrity_if_dag_changed to a bulk UPDATE statement
  • Add session.expire() to ensure relationship cache coherence

Benchmark Results

Benchmarked the actual DagRun.get_running_dag_runs_to_examine() method:

Setup: 20 running DAG runs × 100 tasks = 2,000 TIs (matches DEFAULT_DAGRUNS_TO_EXAMINE)

Version Memory TIs in Memory
NEW (no joinedload) 0.21 MB 0
OLD (with joinedload) 6.56 MB 2,000
Savings 96.9% -

This method is called every scheduler loop (~1/second), so this eliminates ~6.4 MB of memory churn per second.

Run the Benchmark

breeze run pytest airflow-core/tests/unit/jobs/test_benchmark_scheduler_memory.py -xvs
Benchmark Script
"""
Benchmark for scheduler TI loading memory.
Calls the ACTUAL get_running_dag_runs_to_examine() method.
"""
import gc
import tracemalloc

import pytest
from sqlalchemy import select
from sqlalchemy.orm import joinedload

from airflow.models.dagrun import DagRun
from airflow.utils.state import DagRunState, TaskInstanceState


def get_memory_kb():
    current, _ = tracemalloc.get_traced_memory()
    return current / 1024


class TestSchedulerMemoryBenchmark:

    @pytest.fixture
    def setup_running_dag_runs(self, dag_maker, session):
        """Create 20 DAG runs with 100 tasks each."""
        for dag_idx in range(20):
            with dag_maker(dag_id=f"benchmark_dag_{dag_idx}", schedule=None, session=session):
                from airflow.providers.standard.operators.empty import EmptyOperator
                for task_idx in range(100):
                    EmptyOperator(task_id=f"task_{task_idx}")

            dr = dag_maker.create_dagrun(state=DagRunState.RUNNING)
            for ti in dr.task_instances:
                ti.state = TaskInstanceState.SCHEDULED
            session.flush()
        session.commit()
        yield [f"benchmark_dag_{i}" for i in range(20)]

    def test_get_running_dag_runs_to_examine_memory(self, setup_running_dag_runs, session):
        dag_ids = setup_running_dag_runs

        # NEW: Call ACTUAL method (no joinedload)
        session.expunge_all()
        gc.collect()
        tracemalloc.start()
        gc.collect()
        mem_before = get_memory_kb()

        dag_runs_new = list(DagRun.get_running_dag_runs_to_examine(session=session))

        mem_new = get_memory_kb() - mem_before
        tracemalloc.stop()

        print(f"\n[NEW] Memory: {mem_new/1024:.2f} MB, DAG runs: {len(dag_runs_new)}")

        # OLD: Same query WITH joinedload (simulated)
        session.expunge_all()
        gc.collect()
        tracemalloc.start()
        gc.collect()
        mem_before = get_memory_kb()

        dag_runs_old = (
            session.scalars(
                select(DagRun)
                .where(DagRun.state == DagRunState.RUNNING)
                .where(DagRun.dag_id.in_(dag_ids))
                .options(joinedload(DagRun.task_instances))
            )
            .unique()
            .all()
        )
        tis = sum(len(dr.task_instances) for dr in dag_runs_old)

        mem_old = get_memory_kb() - mem_before
        tracemalloc.stop()

        print(f"[OLD] Memory: {mem_old/1024:.2f} MB, TIs loaded: {tis}")
        print(f"SAVINGS: {((mem_old - mem_new) / mem_old) * 100:.1f}%")

        assert mem_new < mem_old

Technical Details

Why both UPDATE and verify_integrity() are needed

  1. Bulk UPDATE: Sets dag_version_id on existing unfinished TIs
  2. verify_integrity(): Creates new TIs for tasks added to DAG

Session synchronization

The session.expire(dag_run, ["task_instances"]) follows SQLAlchemy best practices for bulk operations that bypass the ORM.


Was generative AI tooling used to co-author this PR?
  • Yes - Claude (Opus)

The scheduler's `get_running_dag_runs_to_examine()` was using
`joinedload(DagRun.task_instances)` which loaded ALL TaskInstances
for each DagRun into memory. This could cause significant memory
pressure with large DAGs having many tasks.

Changes:
- Remove unnecessary `joinedload(task_instances)` from the query
- Convert the TI loop in `_verify_integrity_if_dag_changed` to a
  bulk UPDATE statement, avoiding loading TIs into memory entirely
- Add `session.expire()` to ensure relationship cache coherence

This reduces memory usage in the scheduler's hot path, especially
beneficial for DAGs with hundreds of tasks.
@kaxil

kaxil commented Jan 22, 2026

Copy link
Copy Markdown
Member Author

cc @vatsrahul1001 Flagging for testing and review

@kaxil
kaxil requested a review from jedcunningham January 22, 2026 22:04
@jedcunningham
jedcunningham requested a review from Copilot January 22, 2026 22:21

Copilot AI 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.

Pull request overview

This PR addresses a significant memory inefficiency in the Airflow scheduler by removing eager loading of all TaskInstances during DAG run examination, which occurs every scheduler loop (~1/second). The change eliminates unnecessary memory pressure when processing large DAGs.

Changes:

  • Removed joinedload(task_instances) from get_running_dag_runs_to_examine() query
  • Replaced TI iteration loop with bulk UPDATE statement in _verify_integrity_if_dag_changed()
  • Added session cache expiration to maintain relationship coherence after bulk updates

Reviewed changes

Copilot reviewed 2 out of 2 changed files in this pull request and generated 1 comment.

File Description
airflow-core/src/airflow/models/dagrun.py Removes eager loading of task_instances relationship to prevent loading all TIs into memory
airflow-core/src/airflow/jobs/scheduler_job_runner.py Converts TI loop to bulk UPDATE query and adds session expiration for cache coherence

💡 Add Copilot custom instructions for smarter, more guided reviews. Learn how to get started.

Comment thread airflow-core/src/airflow/jobs/scheduler_job_runner.py
@kaxil
kaxil marked this pull request as ready for review January 22, 2026 23:29
@kaxil
kaxil requested review from XD-DENG and ashb as code owners January 22, 2026 23:29
@kaxil
kaxil merged commit dd0eab9 into apache:main Jan 22, 2026
71 checks passed
@kaxil
kaxil deleted the fix-scheduler-ti-loading-memory branch January 22, 2026 23:29
suii2210 pushed a commit to suii2210/airflow that referenced this pull request Jan 26, 2026
…e#60956)

The scheduler's `get_running_dag_runs_to_examine()` was using
`joinedload(DagRun.task_instances)` which loaded ALL TaskInstances
for each DagRun into memory. This could cause significant memory
pressure with large DAGs having many tasks.

Changes:
- Remove unnecessary `joinedload(task_instances)` from the query
- Convert the TI loop in `_verify_integrity_if_dag_changed` to a
  bulk UPDATE statement, avoiding loading TIs into memory entirely
- Add `session.expire()` to ensure relationship cache coherence

This reduces memory usage in the scheduler's hot path, especially
beneficial for DAGs with hundreds of tasks.
shreyas-dev pushed a commit to shreyas-dev/airflow that referenced this pull request Jan 29, 2026
…e#60956)

The scheduler's `get_running_dag_runs_to_examine()` was using
`joinedload(DagRun.task_instances)` which loaded ALL TaskInstances
for each DagRun into memory. This could cause significant memory
pressure with large DAGs having many tasks.

Changes:
- Remove unnecessary `joinedload(task_instances)` from the query
- Convert the TI loop in `_verify_integrity_if_dag_changed` to a
  bulk UPDATE statement, avoiding loading TIs into memory entirely
- Add `session.expire()` to ensure relationship cache coherence

This reduces memory usage in the scheduler's hot path, especially
beneficial for DAGs with hundreds of tasks.
jhgoebbert pushed a commit to jhgoebbert/airflow_Owen-CH-Leung that referenced this pull request Feb 8, 2026
…e#60956)

The scheduler's `get_running_dag_runs_to_examine()` was using
`joinedload(DagRun.task_instances)` which loaded ALL TaskInstances
for each DagRun into memory. This could cause significant memory
pressure with large DAGs having many tasks.

Changes:
- Remove unnecessary `joinedload(task_instances)` from the query
- Convert the TI loop in `_verify_integrity_if_dag_changed` to a
  bulk UPDATE statement, avoiding loading TIs into memory entirely
- Add `session.expire()` to ensure relationship cache coherence

This reduces memory usage in the scheduler's hot path, especially
beneficial for DAGs with hundreds of tasks.
@potiuk

potiuk commented Feb 14, 2026

Copy link
Copy Markdown
Member

Nice!

choo121600 pushed a commit to choo121600/airflow that referenced this pull request Feb 22, 2026
…e#60956)

The scheduler's `get_running_dag_runs_to_examine()` was using
`joinedload(DagRun.task_instances)` which loaded ALL TaskInstances
for each DagRun into memory. This could cause significant memory
pressure with large DAGs having many tasks.

Changes:
- Remove unnecessary `joinedload(task_instances)` from the query
- Convert the TI loop in `_verify_integrity_if_dag_changed` to a
  bulk UPDATE statement, avoiding loading TIs into memory entirely
- Add `session.expire()` to ensure relationship cache coherence

This reduces memory usage in the scheduler's hot path, especially
beneficial for DAGs with hundreds of tasks.
Subham-KRLX pushed a commit to Subham-KRLX/airflow that referenced this pull request Mar 4, 2026
…e#60956)

The scheduler's `get_running_dag_runs_to_examine()` was using
`joinedload(DagRun.task_instances)` which loaded ALL TaskInstances
for each DagRun into memory. This could cause significant memory
pressure with large DAGs having many tasks.

Changes:
- Remove unnecessary `joinedload(task_instances)` from the query
- Convert the TI loop in `_verify_integrity_if_dag_changed` to a
  bulk UPDATE statement, avoiding loading TIs into memory entirely
- Add `session.expire()` to ensure relationship cache coherence

This reduces memory usage in the scheduler's hot path, especially
beneficial for DAGs with hundreds of tasks.
Ankurdeewan pushed a commit to Ankurdeewan/airflow that referenced this pull request Mar 15, 2026
…e#60956)

The scheduler's `get_running_dag_runs_to_examine()` was using
`joinedload(DagRun.task_instances)` which loaded ALL TaskInstances
for each DagRun into memory. This could cause significant memory
pressure with large DAGs having many tasks.

Changes:
- Remove unnecessary `joinedload(task_instances)` from the query
- Convert the TI loop in `_verify_integrity_if_dag_changed` to a
  bulk UPDATE statement, avoiding loading TIs into memory entirely
- Add `session.expire()` to ensure relationship cache coherence

This reduces memory usage in the scheduler's hot path, especially
beneficial for DAGs with hundreds of tasks.
@vatsrahul1001 vatsrahul1001 added this to the Airflow 3.2.0 milestone Apr 7, 2026
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

5 participants