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]

Reply via email to