diff --git a/airflow-core/src/airflow/jobs/scheduler_job_runner.py b/airflow-core/src/airflow/jobs/scheduler_job_runner.py index 089c78f3827e7..6267325dc1118 100644 --- a/airflow-core/src/airflow/jobs/scheduler_job_runner.py +++ b/airflow-core/src/airflow/jobs/scheduler_job_runner.py @@ -1614,7 +1614,7 @@ def process_executor_events( task = dag.get_task(ti.task_id) except Exception: cls.logger().exception("Marking task instance %s as %s", ti, state) - ti.set_state(state) + ti.set_state(state, session=session) continue ti.task = task if task.has_on_retry_callback or task.has_on_failure_callback: @@ -1666,7 +1666,7 @@ def process_executor_events( ) # Adjust max_tries to allow retry beyond normal limits (like clearing does) ti.max_tries = ti.try_number + ti.task.retries - ti.set_state(None) + ti.set_state(None, session=session) continue # Send email notification request to DAG processor via DB diff --git a/airflow-core/tests/unit/jobs/test_scheduler_job.py b/airflow-core/tests/unit/jobs/test_scheduler_job.py index 4d1a540fedacb..e0169bb39d506 100644 --- a/airflow-core/tests/unit/jobs/test_scheduler_job.py +++ b/airflow-core/tests/unit/jobs/test_scheduler_job.py @@ -621,6 +621,45 @@ def test_process_executor_events_restarting_cleared_task(self, mock_task_callbac # Verify try_number wasn't changed (scheduler doesn't increment it here) assert ti1.try_number == 4, "try_number should remain unchanged" + @pytest.mark.parametrize( + ("ti_state", "event_state", "dag_not_found"), + [ + pytest.param(TaskInstanceState.RESTARTING, State.SUCCESS, False, id="cleared-task-terminated"), + pytest.param(TaskInstanceState.QUEUED, State.FAILED, True, id="dag-not-found"), + ], + ) + def test_process_executor_events_sets_state_in_callers_transaction( + self, dag_maker, ti_state, event_state, dag_not_found + ): + """ + Setting a task instance's state must not commit the caller's transaction. + + ``settings.Session`` is scoped, so ``ti.set_state()`` without ``session`` resolved to the + scheduler's own session and ``create_session()`` committed and closed it on exit. That released + the scheduler's row locks mid-batch and detached the task instances still to be processed, so + changes made to them afterwards were never written. + """ + session = settings.Session() + with dag_maker(dag_id="test_executor_events_callers_transaction", fileloc="/test_path1/"): + task1 = EmptyOperator(task_id="test_task", retries=2) + ti1 = dag_maker.create_dagrun().get_task_instance(task1.task_id) + ti1.state = ti_state + session.merge(ti1) + session.commit() + + executor = MockExecutor(do_update=False) + job_runner = SchedulerJobRunner(Job(), executors=[executor]) + if dag_not_found: + job_runner.scheduler_dag_bag = mock.MagicMock() + job_runner.scheduler_dag_bag.get_dag_for_run.side_effect = Exception("failed") + executor.event_buffer[ti1.key] = event_state, None + + job_runner._process_executor_events(executor=executor, session=session) + session.rollback() + + ti1.refresh_from_db(session=session) + assert ti1.state == ti_state + @mock.patch("airflow.jobs.scheduler_job_runner.TaskCallbackRequest") @mock.patch("airflow._shared.observability.metrics.stats._get_backend") def test_process_executor_events_with_no_callback(self, mock_get_backend, mock_task_callback, dag_maker):