Skip to content

Fix short-circuit in a mapped task group not skipping later tasks - #74283

Merged
kaxil merged 3 commits into
apache:mainfrom
astronomer:fix-mapped-group-transitive-skip
Oct 6, 2026
Merged

kaxil merged 3 commits into
apache:mainfrom
astronomer:fix-mapped-group-transitive-skip

Conversation

@kaxil

@kaxil kaxil commented Oct 5, 2026 •

Copy link
Copy Markdown
Member

closes: #74175
related: #73959, #74188

Inside a mapped task group, a short-circuit task that returns False skips only the task right after it. Later tasks of the same map index still run if their trigger rule accepts a skipped upstream (all_done, none_failed), even with the default ignore_downstream_trigger_rules=True. The same chain outside a mapped task group skips them.

The short-circuit already lists every downstream task in its skip decision. For an unmapped task the worker skips those task instances directly. For a mapped one it cannot, because a task inside a mapped task group is only expanded once its own dependencies are met, so the downstream task instances of that map index may not exist yet. SkipMixin leaves the job to NotPreviouslySkippedDep, which only read the decision of direct upstream tasks. #73959 fixed the map index that lookup uses, so the first task is skipped again; this change covers the rest of the list.

Airflow 2.11 has the same gap: its SkipMixin also leaves mapped task instances to NotPreviouslySkippedDep, and that dep also reads direct upstream tasks only. So this is not a regression from 2.x. It makes mapped task groups do what ignore_downstream_trigger_rules=True documents, and what unmapped Dags already do in 2.x and 3.x.

Design rationale

  • Only task instances inside a mapped task group take the new path. It applies to a task instance with a map index whose closest mapped task group holds a SkipMixin task. That task instance is skipped if the group's decision for its map index lists it under skipped. Unmapped Dags keep the existing direct-upstream lookup, with the same queries.
  • One XCom query per scheduling pass. The group's decisions for all map indexes are read once through XComModel.get_many and kept on the DepContext, like the trigger-rule upstream counts. On main, the first task after an in-group short-circuit ran its own query for every map index.
  • followed is still only read for direct downstream tasks. A branch operator lists only its direct children there. Reading it further down would skip every join.
  • Direct parents keep their current rule. Their decision counts once they finish, in any state, which is what main and Airflow 2 do inside and outside mapped task groups. The tasks further down that this change newly reaches need more: the short-circuit must have succeeded. Clearing a task instance keeps its XComs until it runs again, so a short-circuit that was cleared and then skipped or upstream-failed without running would otherwise skip them based on its earlier try.
  • Why not walk all upstream ancestors, as Fix mapped task group short circuit skipping #74188 does? Every task downstream of any SkipMixin task, in every Dag, would then run its own XCom query on every scheduling pass. See the numbers below.

Gotchas

  • Behaviour change. Inside a mapped task group, tasks further down a short-circuit that returned False ran on 2.x and 3.x when their trigger rule accepts a skipped upstream. They are now skipped for that map index, as the ShortCircuitOperator docs describe and as outside a mapped task group. ignore_downstream_trigger_rules=False keeps them running: only the direct downstream is skipped and the rest follow their trigger rules.

Benchmark

One NotPreviouslySkippedDep pass over every schedulable task instance with a shared DepContext, as the scheduler runs it, median of 3 runs in breeze on SQLite. "Mapped group 1000 x 10" is a short-circuit followed by a chain of 10 all_done tasks in a task group expanded 1000 times.

XCom queries per pass:

Scenario Task instances main #74188 this PR
Unmapped, gate then 1999 tasks 1999 1 1999 1
Mapped group 1000 x 10, every gate True 10000 1000 10000 1
Mapped group 1000 x 10, no SkipMixin task 10000 0 0 0
Mapped group 1000 x 10, odd map indexes short-circuited 10000 1000 10000 1

Time, with main and this PR swapped in turn in one container, two rounds:

Scenario main this PR
Unmapped, gate then 1999 tasks 0.148 s, 0.131 s 0.142 s, 0.133 s
Mapped group 1000 x 10, every gate True 2.75 s, 2.84 s 2.57 s, 2.44 s
Mapped group 1000 x 10, no SkipMixin task 2.72 s, 2.52 s 2.49 s, 2.46 s

In a separate run of the same benchmark, #74188 took 0.64 s on the unmapped Dag where main took 0.029 s, and 7.39 s against 2.79 s with every gate True. In the short-circuited scenario main skips 500 task instances (only the first task of each short-circuited map index) and this PR skips all 5000, so most of its time there goes to the extra skip writes. These are SQLite numbers, where a query costs little. On Postgres or MySQL every query is also a network round trip.

Run

airflow standalone (scheduler, LocalExecutor, Postgres) with the Dags triggered through the REST API. Map index 1 has the gate returning False; map index 0 runs every task in these Dags.

Dag group.a[1] group.b[1] group.c[1]
Issue Dag plus a task c, all_done skipped skipped skipped
Same, none_failed skipped skipped skipped
Task group expanded over an upstream task's output, b all_done, c none_failed skipped skipped skipped
ignore_downstream_trigger_rules=False, all_done skipped success success

A task after the group succeeds in every Dag. A branch followed by a join inside a mapped task group skips only the branch not taken, and the join and the task after it succeed for both map indexes. An unmapped short-circuit chain is unchanged. On main, airflow dags test of the first two Dags runs group.b[1] and group.c[1]. The expanded-over-output case also runs with real task execution in test_mappedoperator.py, and fails on main.

Known issues

  • A short-circuit in an outer mapped task group does not skip tasks inside a nested mapped task group. Their map indexes combine both groups, so the outer group's map index does not identify them; that case keeps today's behaviour.
  • Clearing only group.b[1] skips it again, because the short-circuit's decision still lists it. Outside a mapped task group a cleared b would run. Clearing a direct child already behaves this way in both cases.

  • Read the Pull Request Guidelines for more information. Note: commit author/co-author name and email in commits become permanently public when merged.
  • For fundamental code changes, an Airflow Improvement Proposal (AIP) is needed.
  • When adding dependency, check compliance with the ASF 3rd Party License Policy.
  • For significant user-facing changes create newsfragment: {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.

@kaxil
kaxil requested a review from shahar1 October 5, 2026 17:42
@kaxil kaxil added this to the Airflow 3.4.0 milestone Oct 5, 2026
@kaxil
kaxil requested a review from vatsrahul1001 October 5, 2026 17:50
kaxil added 3 commits October 6, 2026 13:52
Inside a mapped task group, a short-circuit task that returned False only skipped
the task right after it. Later tasks of the same map index still ran when their
trigger rule accepts a skipped upstream (all_done, none_failed), although
ignore_downstream_trigger_rules=True lists them in the skip decision.

A task in a mapped task group now honours the skip decision of any SkipMixin task
of the same group for its map index. The decisions are read with one query per
scheduling pass, so tasks outside a mapped task group keep their current lookup.
…more paths

Read the decisions with the same XCom query builder the rest of core uses, and
drop a map index filter that could not change the result.

Add a test that runs a short-circuit inside a task group expanded over an upstream
task's output, where the later tasks are expanded only after the gate has run.
Also check that a gate that failed or was removed after a clear is ignored, and
that one pass evaluating two mapped task groups keeps their decisions apart.
A SkipMixin parent's decision counts once the parent has finished, in any state,
for its direct downstream tasks, as it does outside a mapped task group and in
Airflow 2. Only the tasks further down that this fix newly reaches require the
writer to have succeeded, so a decision left behind by a cleared gate that never
ran again does not skip them.
@kaxil
kaxil force-pushed the fix-mapped-group-transitive-skip branch from 421594d to bd50a26 Compare October 6, 2026 12:52
@kaxil
kaxil merged commit 0ca87cb into apache:main Oct 6, 2026
79 checks passed
@kaxil
kaxil deleted the fix-mapped-group-transitive-skip branch October 6, 2026 17:38
vatsrahul1001 added a commit that referenced this pull request Oct 7, 2026
The mapped-task-group skip-decision query added in #74283 reads
XComModel.task_id/map_index/value directly in with_only_columns, which the
check-xcom-model-columns prek hook forbids, turning the repo-wide static-check
job red for every PR. Read them through xcom_entity(query) as the hook requires.
kaxil pushed a commit that referenced this pull request Oct 7, 2026
…g run (#74390)

* Read XComModel columns through xcom_entity in NotPreviouslySkippedDep

The mapped-task-group skip-decision query added in #74283 reads
XComModel.task_id/map_index/value directly in with_only_columns, which the
check-xcom-model-columns prek hook forbids, turning the repo-wide static-check
job red for every PR. Read them through xcom_entity(query) as the hook requires.

* Fix skip-decision tests for the xcom_v2 rename and add a cross-run guard

The skip-decision query reads xcom_v2 after #74222, so capture_orm_selects
needs that name (the read-once test was silently counting zero). Also add a
regression test that a skip decision from one Dag run does not skip the same
map index in another, which fails on the unscoped query.

* Reformat with ruff format
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Short-circuit in a mapped task group lets later all_done and none_failed tasks run

2 participants