This is an automated email from the ASF dual-hosted git repository.
kaxil pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/airflow.git
The following commit(s) were added to refs/heads/main by this push:
new 884e687ea77 Fix short-circuit inside a mapped task group not skipping
downstream (#73959)
884e687ea77 is described below
commit 884e687ea7736332a3d9ea72bcae6775a9fe689c
Author: Kaxil Naik <[email protected]>
AuthorDate: Thu Oct 1 07:34:57 2026 +0100
Fix short-circuit inside a mapped task group not skipping downstream
(#73959)
---
.../ti_deps/deps/not_previously_skipped_dep.py | 6 ++-
.../deps/test_not_previously_skipped_dep.py | 44 ++++++++++++++++++++++
2 files changed, 48 insertions(+), 2 deletions(-)
diff --git
a/airflow-core/src/airflow/ti_deps/deps/not_previously_skipped_dep.py
b/airflow-core/src/airflow/ti_deps/deps/not_previously_skipped_dep.py
index 4b767a04f23..94417b252fc 100644
--- a/airflow-core/src/airflow/ti_deps/deps/not_previously_skipped_dep.py
+++ b/airflow-core/src/airflow/ti_deps/deps/not_previously_skipped_dep.py
@@ -60,8 +60,10 @@ class NotPreviouslySkippedDep(BaseTIDep):
# Use the parent's map context to look up the XCom. An
unmapped parent
# (e.g. LatestOnlyOperator) writes XCom with map_index=-1, so
we must
- # query with -1 instead of the child's map_index.
- xcom_map_index = ti.map_index if parent.is_mapped else -1
+ # query with -1 instead of the child's map_index. A parent
inside a
+ # mapped task group is expanded like a mapped task, so it
writes XCom
+ # per map index even though ``is_mapped`` is False.
+ xcom_map_index = ti.map_index if parent.get_needs_expansion()
else -1
prev_result = ti.xcom_pull(
task_ids=parent.task_id,
key=XCOM_SKIPMIXIN_KEY,
diff --git
a/airflow-core/tests/unit/ti_deps/deps/test_not_previously_skipped_dep.py
b/airflow-core/tests/unit/ti_deps/deps/test_not_previously_skipped_dep.py
index cd5b364311b..68b0a4f1635 100644
--- a/airflow-core/tests/unit/ti_deps/deps/test_not_previously_skipped_dep.py
+++ b/airflow-core/tests/unit/ti_deps/deps/test_not_previously_skipped_dep.py
@@ -27,11 +27,13 @@ from airflow.models import DagRun, TaskInstance
from airflow.models.xcom import XComModel
from airflow.providers.standard.operators.empty import EmptyOperator
from airflow.providers.standard.operators.python import BranchPythonOperator
+from airflow.sdk import task, task_group
from airflow.sdk.bases.xcom import BaseXCom
from airflow.ti_deps.dep_context import DepContext
from airflow.ti_deps.deps.not_previously_skipped_dep import (
XCOM_SKIPMIXIN_FOLLOWED,
XCOM_SKIPMIXIN_KEY,
+ XCOM_SKIPMIXIN_SKIPPED,
NotPreviouslySkippedDep,
)
from airflow.utils.state import State
@@ -219,6 +221,48 @@ def test_unmapped_parent_skip_mapped_downstream(session,
dag_maker):
assert tis["op2"].state == State.SKIPPED
+def test_parent_in_mapped_task_group_skips_same_map_index(session, dag_maker):
+ """
+ A SkipMixin parent inside a mapped task group writes XCom per map index, so
+ each child TI in the group must read the decision for its own map index.
+ """
+ with dag_maker("test_mapped_group_skip_dag", schedule=None,
session=session):
+
+ @task.short_circuit(task_id="gate")
+ def gate(value):
+ return value
+
+ @task_group
+ def group(value):
+ gate(value) >> EmptyOperator(task_id="child")
+
+ group.expand(value=[True, False])
+
+ dr = dag_maker.create_dagrun(run_type=DagRunType.MANUAL,
state=State.RUNNING)
+ tis = {(ti.task_id, ti.map_index): ti for ti in dr.task_instances}
+ for map_index in (0, 1):
+ tis[("group.gate", map_index)].state = State.SUCCESS
+ session.merge(tis[("group.gate", map_index)])
+ # Only the map index 1 gate short-circuited, as SkipMixin.skip records it.
+ XComModel.set(
+ key=XCOM_SKIPMIXIN_KEY,
+ value={XCOM_SKIPMIXIN_SKIPPED: ["group.child"]},
+ dag_id=dr.dag_id,
+ task_id="group.gate",
+ run_id=dr.run_id,
+ map_index=1,
+ session=session,
+ )
+ session.flush()
+
+ dep = NotPreviouslySkippedDep()
+
+ assert not dep.is_met(tis[("group.child", 1)], session=session)
+ assert tis[("group.child", 1)].state == State.SKIPPED
+ assert dep.is_met(tis[("group.child", 0)], session=session)
+ assert tis[("group.child", 0)].state != State.SKIPPED
+
+
def test_branch_skip_decision_bypasses_custom_xcom_backend(session, dag_maker):
"""
A value-externalizing custom XCom backend must not break branch-skip of