Repository navigation
Conversation
a94b5a9 to
495f5b3
Compare
A serialized Dag that cannot be read is almost always transient -- the Dag processor is mid-write, or the row was briefly unreadable -- but the scheduler treated it as final and failed every SCHEDULED task instance of the Dag, across all of its Dag runs, bypassing retries and failure callbacks. The Dag run was not finished either, so those tasks could only be recovered by clearing them by hand. The bulk update also doubled as a starvation filter: flipping the rows to FAILED is what kept them out of the next iteration of the critical-section query. Starving the run preserves that effect, so an unresolvable Dag cannot hold up the Dags queued behind it. The error is reported once per Dag rather than once per task instance, so a Dag with many stuck runs cannot flood the scheduler log now that its task instances survive the round. closes: apache#62050
495f5b3 to
dd58070
Compare
|
Hello @rjgoyln - thank you for your contributions to Apache Airflow! The Airflow community has introduced a limit of 5 open pull requests at a time for contributors without write access to the repository. You currently have 24 open pull requests, so - as a one-time step of introducing the limit - we closed the ones where maintainers have not engaged yet:
These pull requests stay open because maintainers are already engaged in them - they count towards your limit:
This is not a judgement of you or of your changes. We never told contributors before that opening many pull requests at once was a problem, so there is nothing to feel bad about - and nothing is lost: your branches, commits and the review history stay where they are. What we ask you to do is to make your first prioritization decision: choose which of the pull requests above matter most to you, and reopen them (up to 5 open at a time, including the ones still open) with the "Reopen pull request" button or While your pull requests are waiting for review, the most valuable thing you can do is help in other ways - reviewing other contributors' pull requests, helping with issues, and taking part in the discussions on the devlist and Slack. Why we introduced the limit, what it means for you and how to reopen or restore a pull request is explained in https://github.com/apache/airflow/blob/main/contributing-docs/32_open_pull_request_limit.rst. Drafted-by: Claude Code (Opus 5); reviewed by @potiuk before posting |
Summary
When a run's serialized Dag cannot be resolved during the task-concurrency check, the scheduler currently fails every
SCHEDULEDtask instance of that Dag in one unboundedUPDATE, across all runs.This is problematic for several reasons:
handle_failure, so the affected task instances lose their normal retries, callbacks, and failure logs.RUNNING, leaving manual clearing as the only recovery.dag_versionand the serialized Dag row it resolves to.The
UPDATEalso has a load-bearing side effect: moving the task instances toFAILEDkeeps them out of the next critical-section query. Simply removing the update would allow the same task instances to refill every batch and potentially starve other Dags.Change
Instead of failing the task instances:
SCHEDULEDand skip them for the current scheduler round.The starvation key is
(dag_id, run_id)rather than justdag_idbecause serialized-Dag resolution is performed per run. One run may fail to resolve while another run of the same Dag still resolves successfully.Behavior change
A run whose serialized Dag remains unresolvable now leaves its task instances
SCHEDULEDinstead of marking themFAILED.This preserves the possibility of recovery if the serialized Dag becomes available again. The DagRun cannot progress while its serialized Dag remains unavailable, so the run was already unable to make progress before this change.
Trade-off
If a run never becomes resolvable, its task instances remain
SCHEDULEDand are re-evaluated in subsequent scheduler rounds rather than being removed from consideration by changing their state toFAILED.The additional scheduling cost is limited to cases where these task instances are ahead of work that could otherwise be scheduled. When they fill a batch, the scheduler skips them and continues to the work behind them. They occupy no pool slots and do not count toward task-concurrency limits while they remain
SCHEDULED.There is currently no mechanism that terminates such a run:
dagrun_timeoutis evaluated only after the serialized Dag has been resolved, whiletask_queued_timeoutonly considersQUEUEDtask instances.Handling the terminal lifecycle of a DagRun whose serialized Dag is permanently unavailable is outside the scope of this PR.
Tests
max_tis=2still allow a healthy Dag to queue. The serialized Dag rows are deleted rather than mocked.SCHEDULEDafter a transient miss is queued successfully on the next scheduler round once the serialized Dag becomes available again.Relation to #72652
#72652 addresses the same starvation issue using a per-Dag starvation key.
This change instead tracks starvation per
(dag_id, run_id), deduplicates the resolution error once per Dag per round, and adds regression coverage for the starvation and per-run behavior.closes: #62050
Was generative AI tooling used to co-author this PR?
Generated-by: Claude Code (Opus 5) following the guidelines