This is an automated email from the ASF dual-hosted git repository.
shahar1 pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/airflow.git
The following commit(s) were added to refs/heads/main by this push:
new 7800fc44c07 Fix scheduler crash on restart with queued or running Edge
tasks (#73832)
7800fc44c07 is described below
commit 7800fc44c07a2fc2457a4aee41cf12db9579538c
Author: PoAn Yang <[email protected]>
AuthorDate: Sat Oct 3 01:18:09 2026 +0800
Fix scheduler crash on restart with queued or running Edge tasks (#73832)
* Fix scheduler crash on restart with queued or running Edge tasks
Signed-off-by: PoAn Yang <[email protected]>
* Add the Edge scheduler-crash fix to the edge3 5.0.0 changelog
---------
Signed-off-by: PoAn Yang <[email protected]>
Co-authored-by: Shahar Epstein <[email protected]>
---
providers/edge3/docs/changelog.rst | 1 +
.../providers/edge3/executors/edge_executor.py | 19 +++---
.../unit/edge3/executors/test_edge_executor.py | 71 ++++++++++++++++++++++
3 files changed, 82 insertions(+), 9 deletions(-)
diff --git a/providers/edge3/docs/changelog.rst
b/providers/edge3/docs/changelog.rst
index 1741821c983..65a12ff1dde 100644
--- a/providers/edge3/docs/changelog.rst
+++ b/providers/edge3/docs/changelog.rst
@@ -49,6 +49,7 @@ Features
Bug Fixes
~~~~~~~~~
+* ``Fix scheduler crash on restart with queued or running Edge tasks (#73832)``
* ``Fix airflow edge list-workers always showing null concurrency (#72959)``
Misc
diff --git
a/providers/edge3/src/airflow/providers/edge3/executors/edge_executor.py
b/providers/edge3/src/airflow/providers/edge3/executors/edge_executor.py
index 6990b0097b7..cbc5e5e9811 100644
--- a/providers/edge3/src/airflow/providers/edge3/executors/edge_executor.py
+++ b/providers/edge3/src/airflow/providers/edge3/executors/edge_executor.py
@@ -43,7 +43,7 @@ from airflow.providers.edge3.models.types import (
from airflow.providers.edge3.version_compat import AIRFLOW_V_3_4_PLUS
from airflow.utils.db import DBLocks, create_global_lock
from airflow.utils.helpers import prune_dict
-from airflow.utils.session import NEW_SESSION, provide_session
+from airflow.utils.session import NEW_SESSION, create_session, provide_session
from airflow.utils.state import TaskInstanceState
if AIRFLOW_V_3_4_PLUS:
@@ -422,10 +422,7 @@ class EdgeExecutor(BaseExecutor):
)
self.log.info("Revoked task instance %s from EdgeExecutor", ti.key)
- @provide_session
- def try_adopt_task_instances(
- self, tis: Sequence[TaskInstance], *, session: Session = NEW_SESSION
- ) -> Sequence[TaskInstance]:
+ def try_adopt_task_instances(self, tis: Sequence[TaskInstance]) ->
Sequence[TaskInstance]:
"""
Adopt the task instances whose job is still in flight in the edge_job
table.
@@ -435,10 +432,14 @@ class EdgeExecutor(BaseExecutor):
:return: any TaskInstances that were unable to be adopted
"""
- tracked_keys = self._get_tracked_job_keys(
- session,
- states=(TaskInstanceState.QUEUED, TaskInstanceState.RESTARTING,
TaskInstanceState.RUNNING),
- )
+ # The scheduler calls this without passing a session, while its own
scoped session still holds
+ # ``tis``. A scoped session here would be that same session, and
create_session() would commit
+ # and close it on exit, which detaches ``tis`` before the scheduler
reads them again.
+ with create_session(scoped=False) as session:
+ tracked_keys = self._get_tracked_job_keys(
+ session,
+ states=(TaskInstanceState.QUEUED,
TaskInstanceState.RESTARTING, TaskInstanceState.RUNNING),
+ )
self.running.update(ti.key for ti in tis if ti.key in tracked_keys)
return [ti for ti in tis if ti.key not in tracked_keys]
diff --git a/providers/edge3/tests/unit/edge3/executors/test_edge_executor.py
b/providers/edge3/tests/unit/edge3/executors/test_edge_executor.py
index 463c515cc3f..bb636de8382 100644
--- a/providers/edge3/tests/unit/edge3/executors/test_edge_executor.py
+++ b/providers/edge3/tests/unit/edge3/executors/test_edge_executor.py
@@ -29,6 +29,8 @@ import time_machine
from sqlalchemy import delete, select
from airflow.executors.workloads import BundleInfo, ExecuteTask
+from airflow.jobs.job import Job
+from airflow.jobs.scheduler_job_runner import SchedulerJobRunner
from airflow.models.taskinstance import TaskInstance
from airflow.providers.common.compat.sdk import Stats, TaskInstanceKey, conf,
timezone
from airflow.providers.edge3.executors.edge_executor import EdgeExecutor
@@ -38,6 +40,7 @@ from airflow.providers.edge3.models.types import
EXECUTE_CALLBACK_TAG
from airflow.utils.session import create_session
from airflow.utils.state import TaskInstanceState
+from tests_common.test_utils.compat import EmptyOperator
from tests_common.test_utils.config import conf_vars
from tests_common.test_utils.version_compat import AIRFLOW_V_3_2_PLUS,
AIRFLOW_V_3_3_PLUS
@@ -407,6 +410,74 @@ class TestEdgeExecutor:
assert key not in executor.running
assert key not in executor.queued_tasks
+ def test_try_adopt_task_instances_keeps_caller_session_open(self):
+ key = TaskInstanceKey(
+ dag_id="test_dag", run_id="test_run", task_id="test_task",
map_index=-1, try_number=1
+ )
+ with create_session() as session:
+ session.add(
+ EdgeJobModel(
+ dag_id="test_dag",
+ task_id="test_task",
+ run_id="test_run",
+ map_index=-1,
+ try_number=1,
+ state=TaskInstanceState.QUEUED,
+ queue="default",
+ command="mock",
+ concurrency_slots=1,
+ )
+ )
+ session.commit()
+ executor = EdgeExecutor()
+
+ # The scheduler calls try_adopt_task_instances() without a session and
keeps using the objects
+ # it loaded into its own scoped session.
+ with create_session() as session:
+ job = session.scalar(select(EdgeJobModel))
+ executor.try_adopt_task_instances([mock.Mock(spec=TaskInstance,
key=key)])
+
+ assert job in session
+ assert executor.running == {key}
+
+
@mock.patch("airflow.executors.executor_loader.ExecutorLoader.init_executors",
autospec=True)
+ def test_scheduler_restart_adopts_queued_edge_task(self,
mock_init_executors, dag_maker, session):
+ with dag_maker("test_dag", session=session):
+ EmptyOperator(task_id="test_task")
+ dag_run = dag_maker.create_dagrun()
+ previous_scheduler_job = Job()
+ restarted_scheduler_job = Job()
+ session.add_all([previous_scheduler_job, restarted_scheduler_job])
+ session.flush()
+ ti = dag_run.get_task_instance("test_task", session=session)
+ ti.state = TaskInstanceState.QUEUED
+ ti.queued_by_job_id = previous_scheduler_job.id
+ session.add(
+ EdgeJobModel(
+ dag_id=ti.dag_id,
+ task_id=ti.task_id,
+ run_id=ti.run_id,
+ map_index=ti.map_index,
+ try_number=ti.try_number,
+ state=TaskInstanceState.QUEUED,
+ queue="default",
+ command="mock",
+ concurrency_slots=1,
+ )
+ )
+ session.commit()
+ # A restarted scheduler loads the task instance from the database, not
from this session.
+ session.expunge_all()
+ executor = EdgeExecutor()
+ mock_init_executors.return_value = [executor]
+
+ SchedulerJobRunner(job=restarted_scheduler_job,
num_runs=0).adopt_or_reset_orphaned_tasks()
+
+ ti.refresh_from_db(session=session)
+ assert ti.state == TaskInstanceState.QUEUED
+ assert ti.queued_by_job_id == restarted_scheduler_job.id
+ assert executor.running == {ti.key}
+
class TestEdgeExecutorMultiTeam:
"""Tests for multi-team (AIP-67) support in EdgeExecutor."""