shahar1 commented on code in PR #74188:
URL: https://github.com/apache/airflow/pull/74188#discussion_r4178451878


##########
airflow-core/src/airflow/ti_deps/deps/not_previously_skipped_dep.py:
##########
@@ -46,7 +46,7 @@ class NotPreviouslySkippedDep(BaseTIDep):
     def _get_dep_statuses(self, ti, dep_context, *, session):
         from airflow.utils.state import TaskInstanceState
 
-        upstream = ti.task.get_direct_relatives(upstream=True)
+        upstream = ti.task.get_flat_relatives(upstream=True)

Review Comment:
   **Blocking:** this regresses every branch + join Dag.
   
   `SkipMixin.skip_all_except` stores `{"followed": [...]}` with only the 
**direct** children it followed. The `followed` rule below ("skip anything not 
in `followed`") now fires for every transitive descendant of a branch operator, 
so a `join` downstream of `task1` gets `SKIPPED`:
   
   ```text
   branch --> task1 --> join   (none_failed_min_one_success)
      \-----> task2 --/
   ```
   
   Verified locally on this branch: `get_dep_statuses(join)` returns 
`passed=False, reason='Skipping because of previous XCom result from parent 
task branch'`; on `main` it returns `[]`.
   
   Keep the flat walk, but only apply the `followed` rule to direct parents 
(the explicit `skipped` list is safe from any ancestor):
   
   ```suggestion
           direct_upstream_ids = {t.task_id for t in 
ti.task.get_direct_relatives(upstream=True)}
           upstream = ti.task.get_flat_relatives(upstream=True)
   ```
   
   and guard the `followed` branch with `parent.task_id in direct_upstream_ids 
and ...`. With that, your new test, the existing suite and a branch + join 
probe all pass (10/10). Full probe test in the review body; please add it to 
this file.
   
   Also worth noting in the description: `get_flat_relatives` walks all 
ancestors on every dep evaluation, which is more scheduler work than the 
previous direct lookup on large Dags.



##########
airflow-core/tests/unit/ti_deps/deps/test_not_previously_skipped_dep.py:
##########
@@ -262,6 +263,56 @@ def group(value):
     assert dep.is_met(tis[("group.child", 0)], session=session)
     assert tis[("group.child", 0)].state != State.SKIPPED
 
+def test_parent_in_mapped_task_group_skips_transitive_downstream(session, 
dag_maker):

Review Comment:
   nit: `ruff format` wants two blank lines before this `def` (and reflows the 
`>>` chain below). The static-checks job will fail as-is; `prek run ruff-format 
--from-ref main` fixes it.



-- 
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