Fix scheduler crash when resetting orphaned task instances - #71348
akhilpratap1991 wants to merge 2 commits into
Conversation
|
Congratulations on your first Pull Request and welcome to the Apache Airflow community! If you have any issues or are unsure about any anything please check our Contributors' Guide
|
|
Static checks need fixing. The description asserts this was "observed in production with KubernetesExecutor as a scheduler CrashLoopBackOff," but the test only reproduces it by manually calling |
|
Static checks are green on the current head — the initial run was sitting in the first-time-contributor approval gate ( On the real code path — the concrete trigger is an executor opening and closing a scoped session while adopting:
#67850 already fixed that one call site provider-side ( I've updated the test accordingly: the executor mock now performs the exact operation the released provider does ( Drafted-by: Claude Code (Opus 4.8); reviewed by @akhilpratap1991 before posting |
Executors' try_adopt_task_instances() can hand back TaskInstances that are no longer bound to the scheduler's session, while the orphan query loads only a handful of columns. The reset path then copies the full row into TaskInstanceHistory and the adopt path reads last_heartbeat_at and dag_run.conf — reads that trigger lazy loads, which on a detached instance raise DetachedInstanceError. The exception is unhandled in _run_scheduler_loop, so the scheduler exits; because it dies before the orphan is reset, the same orphan crashes it again on every restart — an unrecoverable crash loop that halts all scheduling (observed in production with KubernetesExecutor as a CrashLoopBackOff). Where the attributes happen to be loaded, the writes to the detached instances are silently lost instead, so orphans are never actually reset. Re-selecting the rows by their already-loaded ids gives both branches fully loaded, session-bound instances, making the outcome independent of what the executor did with the objects it was handed.
Review asked for the concrete code path that detaches the TaskInstances rather than simulating it with session.expunge(). The mechanism is the executor opening and closing a scoped session while adopting: create_session() returns the scheduler's own thread-scoped session and closes it on exit, expunging everything the orphan query loaded. KubernetesExecutor's completed-pod adoption did exactly this from apache#66400 until apache#67850, and released providers up to 10.17.x still ship it (the combination Airflow 3.2.2 constraints pin). The mock now performs that exact operation, so the test exercises the production mechanism end-to-end instead of hand-detaching instances.
ff338a9 to
13fb904
Compare
adopt_or_reset_orphaned_tasks()reads deferred columns (the full row copied intoTaskInstanceHistorybyprepare_db_for_next_try,last_heartbeat_at) and the lazydag_runrelationship on TaskInstances that executors'try_adopt_task_instances()may hand back detached from the session. Those reads raiseDetachedInstanceError, which is unhandled in_run_scheduler_loop— the scheduler exits, and since it dies before the orphan is reset, the same orphan crashes it again on every restart: an unrecoverable crash loop that halts all scheduling (observed in production with KubernetesExecutor as a schedulerCrashLoopBackOff). Where the attributes happen to be loaded, the writes to the detached instances are silently lost instead, so the orphans are never actually reset.closes: #71272
This completes #67822, which fixed the first such read (
repr(ti)onstate/map_index) by wideningload_only— butTaskInstanceHistory(ti)copies ~35 columns, so widening the column list further just moves the crash. Instead, re-select the rows to mutate by their already-loaded ids (withdag_runjoined-loaded) so both the reset and adopt branches operate on fully loaded, session-bound instances, independent of what the executor did with the objects it was handed.The regression test simulates the executor detaching the TaskInstances it was handed (as observed in production): without this change it fails with the production
DetachedInstanceError(deferred load oftry_number); with it, the orphan is reset with aTaskInstanceHistoryaudit row and the adopted TI gets its heartbeat/dag_run handling.Was generative AI tooling used to co-author this PR?
Generated-by: Claude Code (Opus 4.8) following the guidelines