Repository navigation
Make EdgeExecutor respect [core] parallelism - #72048
Conversation
dheerajturaga
left a comment
There was a problem hiding this comment.
Requesting changes for two correctness issues:
-
The worker claim transition can incorrectly fail a task.
queue_workload()now adds the task key toself.running, while the worker's fetch endpoint commits the job asRESTARTINGbefore it reportsRUNNING._purge_jobs()treats everyRESTARTINGjob inself.runningas failed. A scheduler heartbeat in that window therefore emits aFAILEDexecutor event and removes the slot even though the worker is starting the task. This needs a non-terminal claim state or equivalent handling that distinguishes a worker claim from an actual failed/retry transition. -
Scheduler failover loses the new parallelism accounting.
try_adopt_task_instances()returns an empty list, declaring all Edge tasks adopted, but it never restores their keys to the new executor'sself.runningset._get_tracked_job_keys()only intersects the existing set, so it cannot restore those entries. After a scheduler restart, existing Edge tasks consume no executor slots and the scheduler can queue up toparallelismadditional tasks. Adoption should populateself.runningfor the Edge jobs actually present for the executor's team and return any missing jobs as not adopted.
I reproduced both cases with focused Edge provider tests: the claim-state test received a FAILED event, and the adoption test left running empty with the slot still available.
Drafted-by: Codex (GPT-5); reviewed by @dheerajturaga before posting
e97aa54 to
d3a44bf
Compare
|
@dheerajturaga Good catch! Thanks for the review. I did following changes.
|
dheerajturaga
left a comment
There was a problem hiding this comment.
Thanks for the quick turnaround on the previous round — both fixes look correct: the RESTARTING-vs-failure change in _purge_jobs() matches the fetch endpoint's claim semantics, and try_adopt_task_instances() now restores running from the team's edge_job rows.
Going through the rest of the diff turned up a few more correctness issues, most to least severe:
-
Callback workloads bypass the new parallelism enforcement (
edge_executor.py:106) — the callback-execute branch ofqueue_workload()(lines 106-134) never callsself.running.add(...), unlike theExecuteTaskbranch (line 171). So DAG/task callback workloads still bypass the parallelism enforcement this PR introduces — an unbounded number can run concurrently regardless of[core] parallelism, even thoughslots_available's docstring says it should count "tasks, callbacks, and connection tests." This is the one I'd want addressed before merge, since it undermines the PR's own goal for that workload type. -
try_adopt_task_instances()can adopt already-finished jobs (edge_executor.py:419) —_get_tracked_job_keys()(lines 264-282) has no state filter, so it returns keys forSUCCESS/FAILED/REMOVEDrows too, as long as they haven't yet been purged by_purge_jobs()'s success/fail purge timers. After a scheduler restart, this can re-adopt an already-finished job and hold a parallelism slot for it until the next purge cycle. -
self.runningupdated before the caller's session commit (edge_executor.py:171) —self.running.add(key)runs before the caller (e.g. the scheduler's_critical_section_enqueue_task_instances, which batches multiple TIs per commit) commits the session. If that transaction rolls back, the DB insert is undone but the in-memory key isn't —self.runningover-reports occupied slots until the next_purge_jobs()sync intersects it with the DB. Self-healing within one heartbeat interval, so lower severity than the other two.
Apologies for the back-and-forth here — appreciate you working through these with me.
Drafted-by: Claude Code (Sonnet 5); reviewed by @dheerajturaga before posting
|
Thanks for the second review. Learn a lot from you comment.
|
93bf09a to
ec4969d
Compare
dheerajturaga
left a comment
There was a problem hiding this comment.
Thanks for addressing the earlier findings. One blocking regression remains:
build_job_key() treats every row with dag_id == "ExecuteCallback" as a callback. Since ExecuteCallback is a valid Dag ID, normal tasks in such a Dag are reconstructed as CallbackKeys, removed from self.running during synchronization, and no longer report terminal states.
Please identify callbacks using the complete synthetic identity—such as the expected run_id, task_id, try_number, and map_index combination—or persist an explicit workload type. Please also add a regression test for a normal task in a Dag named ExecuteCallback.
Once the key collision is fixed, this looks ready to merge.
Drafted-by: Codex (GPT-5); reviewed by @dheerajturaga before posting
Signed-off-by: PoAn Yang <payang@apache.org>
Signed-off-by: PoAn Yang <payang@apache.org>
Signed-off-by: PoAn Yang <payang@apache.org>
Signed-off-by: PoAn Yang <payang@apache.org>
ec4969d to
9f61bbe
Compare
|
Thanks for the review. Besides the key collision, this push fixes one more problem and add a changelog about this PR:
|
dheerajturaga
left a comment
There was a problem hiding this comment.
Great! thanks for implementing this and going back and forth with reviews.
Why
BaseExecutorworkload pipeline and record each workload inself.running.KubernetesExecutordoes it in_process_workloads(), soslots_availablereflects what they have in flight.EdgeExecutorwrites queued tasks straight into theedge_jobtable and never records them inself.running. Nothing has written to that set since Remove Airflow 2 code path in executors #51009 removed the Airflow 2 dispatch path.slots_availableis therefore alwaysparallelismand[core] parallelismhas no effect onEdgeExecutor.if job.key in self.running:block in_purge_jobs()unreachable, so the executor never reports task state to the scheduler.How
queue_workload()records the key of anExecuteTaskworkload inself.running._purge_jobs()reconcilesself.runningagainst every job row of the team, queued ones included. The previous reconciliation covered only the six non-queued states, so a key added at queue time was dropped again on the nextsync(), before a worker could pick the job up. The_get_tracked_job_keysread takes no row locks: an edge worker fetches its next job withFOR UPDATE SKIP LOCKED, so locking queued rows here would make it come back empty._purge_jobs()no longer reportsRESTARTINGas a failure. The fetch endpoint parks a claimed job in that state until the worker reportsRUNNING, so the executor now records it as a transition and the slot stays taken.EdgeJobModel.keyreturns theairflow.modelsTaskInstanceKeyinstead of theairflow.sdkone.state_class_for_key()dispatches withisinstanceand the two are unrelated classes, sorunning_state()would raiseTypeErroras soon asrunningbecame non-empty.try_adopt_task_instances()restoresrunningfrom the team'sedge_jobrows after a scheduler restart and returns task instances without a job row as not adopted.Verification
uv run --project providers/edge3 pytest providers/edge3/tests/unitWas generative AI tooling used to co-author this PR?
{pr_number}.significant.rst, in airflow-core/newsfragments. You can add this file in a follow-up commit after the PR is created so you know the PR number.