kaxil opened a new pull request, #74283: URL: https://github.com/apache/airflow/pull/74283
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. ## Design rationale - **Only task instances inside a mapped task group take the new path.** A task instance with a map index whose closest mapped task group holds a `SkipMixin` task reads that group's decisions for its map index, and is skipped if the decision lists it under `skipped`. Every other task instance runs the existing direct-upstream lookup, so unmapped Dags see no change in behaviour or in queries. - **One XCom query per scheduling pass instead of one per task instance.** The decisions for every map index of the group are read once and kept on the `DepContext`, the same way the trigger-rule upstream counts are kept there. This also drops the query the first task after an in-group short-circuit ran per map index on main. - **`followed` still applies only to direct downstream tasks.** A branch operator's `followed` list names its direct children only, so honouring it further down would skip every join. - **A decision counts only if the task instance that wrote it succeeded.** Clearing a task instance keeps its XComs until it runs again. Without this, a short-circuit that was cleared and then skipped or upstream-failed without running would skip tasks based on its earlier try. - **Why not walk all upstream ancestors, as #74188 does?** That applies to every task in every Dag: each task downstream of any `SkipMixin` task then runs its own XCom query on every scheduling pass, mapped or not. The numbers are below. ## 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. | Scenario | Task instances | XCom queries (main / #74188 / this PR) | Time (main / #74188 / this PR) | |---|---|---|---| | Unmapped, gate then 1999 tasks | 1999 | 1 / 1999 / 1 | 0.029 s / 0.64 s / 0.026 s | | Mapped group 1000 x 10, every gate True | 10000 | 1000 / 10000 / **1** | 2.79 s / 7.39 s / 2.34 s | | Mapped group 1000 x 10, no `SkipMixin` task | 10000 | 0 / 0 / 0 | 2.54 s / 3.05 s / 2.40 s | | Mapped group 1000 x 10, odd map indexes short-circuited | 10000 | 1000 / 10000 / 1 | 3.84 s / 15.65 s / 10.16 s | In the last row main skips 500 task instances (only the first task of each short-circuited map index) and the other two skip all 5000. Most of the time there is the 5000 skip writes. SQLite queries are cheap; on Postgres or MySQL each saved query is a round trip. ## Run `airflow dags test` on the Dag from the issue, extended with a third task `c` and a task `after` downstream of the group, map index 1 (the gate returns False). Map index 0 runs everything in every case. | Dag | `group.a[1]` | `group.b[1]` | `group.c[1]` | |---|---|---|---| | `all_done`, main | skipped | success | success | | `none_failed`, main | skipped | success | success | | `all_done`, this PR | skipped | skipped | skipped | | `none_failed`, this PR | skipped | skipped | skipped | | `ignore_downstream_trigger_rules=False`, this PR | skipped | success | success | With this PR `after` succeeds in every case. 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. ## 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](https://github.com/apache/airflow/blob/main/contributing-docs/05_pull_requests.rst#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](https://cwiki.apache.org/confluence/display/AIRFLOW/Airflow+Improvement+Proposals)) is needed. * When adding dependency, check compliance with the [ASF 3rd Party License Policy](https://www.apache.org/legal/resolved.html#category-x). * For significant user-facing changes create newsfragment: `{pr_number}.significant.rst`, in [airflow-core/newsfragments](https://github.com/apache/airflow/tree/main/airflow-core/newsfragments). You can add this file in a follow-up commit after the PR is created so you know the PR number. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
