aaron-y-chen commented on code in PR #73470:
URL: https://github.com/apache/airflow/pull/73470#discussion_r4088796692
##########
task-sdk/src/airflow/sdk/bases/skipmixin.py:
##########
@@ -79,7 +79,14 @@ def skip(
value={XCOM_SKIPMIXIN_SKIPPED: task_ids_list},
)
- self._set_state_to_skipped(task_ids_list, ti.map_index)
+ # A task with wait_for_past_depends_before_skipping=True must not be
force-skipped
+ # here: doing so bypasses NotPreviouslySkippedDep, the only place that
flag is
+ # honored, so it never gets the chance to hold the task back until its
own past
+ # depends are met. Leave it for that dependency check to decide
instead.
+ immediate_skip_ids = [
+ d.task_id for d in task_list if not getattr(d,
"wait_for_past_depends_before_skipping", False)
Review Comment:
I reproduced a regression in the default `ShortCircuitOperator` mode. With
`gate >> middle >> target`, where `target` has `depends_on_past=True`,
`wait_for_past_depends_before_skipping=True`, and `trigger_rule="all_done"`,
run 1 succeeds and `gate` returns False in run 2.
On this PR, `middle` is skipped but the scheduler considers `target` ready
and marks it successful. On the base commit, `target` is skipped. `skip()`
records `target` in gate's XCom, but `NotPreviouslySkippedDep` only checks
direct upstream tasks, so it cannot find that decision through `middle`. Could
we preserve the default behavior of skipping all descendants when deferring
this skip, and add a two-run regression test?
--
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]