ashb commented on code in PR #74222:
URL: https://github.com/apache/airflow/pull/74222#discussion_r4193958316


##########
airflow-core/src/airflow/models/xcom.py:
##########
@@ -399,6 +424,201 @@ def deserialize_value(result: Any) -> Any:
             return result.value
 
 
+def _rows():
+    from airflow.models.dagrun import DagRun
+    from airflow.models.taskinstance import LegacyTaskDataOwner, TaskInstance
+
+    coordinates = ("dag_id", "task_id", "run_id", "map_index")
+    data = ("key", "value", "timestamp", "dag_result", "mapped_length")
+    run_join = and_(DagRun.dag_id == TaskInstance.dag_id, DagRun.run_id == 
TaskInstance.run_id)
+    context = [
+        *(getattr(TaskInstance, name) for name in coordinates),
+        DagRun.id.label("dag_run_id"),
+        DagRun.logical_date,
+        DagRun.run_after,
+    ]
+    new_rows = (
+        select(XComModelV2.task_instance_id, *(getattr(XComModelV2, name) for 
name in data), *context)
+        .select_from(XComModelV2)
+        .join(TaskInstance, XComModelV2.task_instance_id == TaskInstance.id)
+        .join(DagRun, run_join)
+    )
+    old_rows = (
+        select(LegacyTaskDataOwner.task_instance_id, *(getattr(XComModelV1, 
name) for name in data), *context)
+        .select_from(LegacyTaskDataOwner)
+        .join(
+            XComModelV1,
+            and_(*(getattr(LegacyTaskDataOwner, name) == getattr(XComModelV1, 
name) for name in coordinates)),
+        )
+        .join(TaskInstance, LegacyTaskDataOwner.task_instance_id == 
TaskInstance.id)
+        .join(DagRun, run_join)
+        .where(
+            ~select(1)
+            .select_from(XComModelV2)
+            .where(
+                XComModelV2.task_instance_id == 
LegacyTaskDataOwner.task_instance_id,
+                XComModelV2.key == XComModelV1.key,
+            )
+            .exists()
+        )
+    )
+    return union_all(new_rows, old_rows)
+
+
+def _filter_rows(rows: Subquery, *, producer_ids: Select, key: str | None) -> 
Subquery:

Review Comment:
   This should be tested by behaviour, rather than explicit SQL generation 
test, in test_xcoms.py::test_retired_attempt_read_response_by_version and 
test_retired_attempt_cannot_replace_or_delete_xcom.



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