diff --git a/airflow-core/src/airflow/jobs/scheduler_job_runner.py b/airflow-core/src/airflow/jobs/scheduler_job_runner.py index 3e9e2321bc9c3..813f782437a04 100644 --- a/airflow-core/src/airflow/jobs/scheduler_job_runner.py +++ b/airflow-core/src/airflow/jobs/scheduler_job_runner.py @@ -1869,8 +1869,10 @@ def _do_scheduling(self, session: Session) -> int: guard.commit() # Bulk fetch the currently active dag runs for the dags we are - # examining, rather than making one query per DagRun - dag_runs = DagRun.get_running_dag_runs_to_examine(session=session) + # examining, rather than making one query per DagRun. + # Materialize into a list because the multi-team block below iterates + # the result and ScalarResult is a one-pass iterator. + dag_runs = list(DagRun.get_running_dag_runs_to_examine(session=session)) if self._multi_team and dag_runs: unique_dag_ids = {dr.dag_id for dr in dag_runs} diff --git a/airflow-core/tests/unit/jobs/test_scheduler_job.py b/airflow-core/tests/unit/jobs/test_scheduler_job.py index 3ee51bb917dc3..1a0897806331a 100644 --- a/airflow-core/tests/unit/jobs/test_scheduler_job.py +++ b/airflow-core/tests/unit/jobs/test_scheduler_job.py @@ -9343,6 +9343,43 @@ def test_multi_team_sets_team_name_on_task_instances(self, dag_maker, mock_execu assert ti._team_name == "team_a" assert ti.stats_tags == {"dag_id": "dag_a", "task_id": "task_a", "team_name": "team_a"} + @conf_vars({("core", "multi_team"): "true"}) + def test_do_scheduling_multi_team_schedules_task_instances(self, dag_maker, session): + """Test that _do_scheduling correctly schedules tasks when multi_team is enabled. + + Regression test: the multi-team code path used to consume the ScalarResult iterator + (returned by get_running_dag_runs_to_examine) when building the team-name mapping, + leaving an exhausted iterator for _schedule_all_dag_runs. This caused tasks to remain + in None state indefinitely. + """ + clear_db_teams() + clear_db_dag_bundles() + + team = Team(name="team_a") + session.add(team) + session.flush() + + bundle = DagBundleModel(name="bundle_a") + bundle.teams.append(team) + session.add(bundle) + session.flush() + + with dag_maker(dag_id="test_multi_team_scheduling", bundle_name="bundle_a", session=session): + EmptyOperator(task_id="task1") + + scheduler_job = Job() + self.job_runner = SchedulerJobRunner(job=scheduler_job, executors=[self.null_exec]) + + dr = dag_maker.create_dagrun(state=State.RUNNING) + ti = dr.get_task_instance("task1", session=session) + assert ti.state == State.NONE + + self.job_runner._do_scheduling(session) + + ti = session.merge(ti) + session.refresh(ti) + assert ti.state != State.NONE + @pytest.mark.need_serialized_dag def test_schedule_dag_run_with_upstream_skip(dag_maker, session):