This is an automated email from the ASF dual-hosted git repository.
potiuk 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 2288ea4537a Fix XCom query in not_previously_skipped_dep to use
xcom_entity (#74377)
2288ea4537a is described below
commit 2288ea4537a4f3e93de6f47a7461d82d9a8a0998
Author: Aaron Chen <[email protected]>
AuthorDate: Tue Oct 6 23:52:39 2026 -0700
Fix XCom query in not_previously_skipped_dep to use xcom_entity (#74377)
* Fix XCom query in not_previously_skipped_dep to use xcom_entity
* Fix XCom capture in mapped task group tests to use xcom_v2
---
airflow-core/src/airflow/ti_deps/deps/not_previously_skipped_dep.py | 5 +++--
.../tests/unit/ti_deps/deps/test_not_previously_skipped_dep.py | 4 ++--
2 files changed, 5 insertions(+), 4 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 1af358b07ed..16c21d8308a 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
@@ -20,7 +20,7 @@ from __future__ import annotations
from typing import TYPE_CHECKING
from airflow.models.taskinstance import PAST_DEPENDS_MET
-from airflow.models.xcom import XComModel
+from airflow.models.xcom import XComModel, xcom_entity
from airflow.ti_deps.deps.base_ti_dep import BaseTIDep
from airflow.utils.state import TaskInstanceState
@@ -177,8 +177,9 @@ def _mapped_group_skip_decisions(
query = XComModel.get_many(
run_id=ti.run_id, key=XCOM_SKIPMIXIN_KEY, dag_ids=ti.dag_id,
task_ids=skipmixin_task_ids
)
+ entity = xcom_entity(query)
rows = session.execute(
- query.with_only_columns(XComModel.task_id, XComModel.map_index,
XComModel.value).order_by(None)
+ query.with_only_columns(entity.task_id, entity.map_index,
entity.value).order_by(None)
)
for row in rows:
if (state := finished_states.get((row.task_id, row.map_index))) is
None:
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 180ed246f31..4a17e26632e 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
@@ -537,7 +537,7 @@ def
test_mapped_task_group_skip_decisions_read_once_per_pass(session, dag_maker)
dep_context =
DepContext(finished_tis=dr.get_task_instances(state=State.finished,
session=session))
dep = NotPreviouslySkippedDep()
- with capture_orm_selects("xcom") as statements:
+ with capture_orm_selects("xcom_v2") as statements:
met = {
(task_id, map_index): dep.is_met(tis[(task_id, map_index)],
dep_context, session=session)
for task_id in downstream
@@ -565,7 +565,7 @@ def
test_mapped_task_group_without_skipmixin_reads_no_xcom(session, dag_maker):
dep_context =
DepContext(finished_tis=dr.get_task_instances(state=State.finished,
session=session))
dep = NotPreviouslySkippedDep()
- with capture_orm_selects("xcom") as statements:
+ with capture_orm_selects("xcom_v2") as statements:
assert all(dep.is_met(tis[("group.b", i)], dep_context,
session=session) for i in range(3))
assert statements == []