kaxil commented on code in PR #73315:
URL: https://github.com/apache/airflow/pull/73315#discussion_r4201398308


##########
airflow-core/src/airflow/serialization/definitions/operatorlink.py:
##########
@@ -55,17 +86,10 @@ def get_link(self, operator: Operator, *, ti_key: 
TaskInstanceKey) -> str:
             "Attempting to retrieve link from XComs with key: %s for task id: 
%s", self.xcom_key, ti_key
         )
         with create_session() as session:
-            result = session.execute(
-                XComModel.get_many(
-                    key=self.xcom_key,
-                    run_id=ti_key.run_id,
-                    dag_ids=ti_key.dag_id,
-                    task_ids=ti_key.task_id,
-                    map_indexes=ti_key.map_index,
-                )
-                .with_only_columns(XComModel.value)
-                .limit(1)
-            ).first()
+            # Runs from before per-try keys existed only have the unsuffixed 
key.
+            result = self._read_value(

Review Comment:
   Since #74222 landed on main, I think most of this is already handled there. 
Each retry now runs under a fresh TI UUID and `ti_run` only clears XComs that 
the current attempt produced, so an earlier try's link rows aren't touched. On 
the read side, `XComOperatorLink.get_link` already filters on 
`try_number=ti_key.try_number`.
   
   Is anything left here once you rebase? This branch conflicts with main now. 
The bit I couldn't work out is the unsuffixed fallback for rows written before 
the upgrade; that would depend on how #74222 maps `xcom_v1` rows to a try 
number.



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