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


##########
airflow-core/src/airflow/serialization/definitions/operatorlink.py:
##########
@@ -43,6 +53,20 @@ class XComOperatorLink(LoggingMixin):
     name: str
     xcom_key: str
 
+    @staticmethod
+    def _read_value(session: Session, key: str, ti_key: TaskInstanceKey) -> 
Row | None:
+        return session.execute(
+            XComModel.get_many(

Review Comment:
   [comements from 
ash](https://github.com/apache/airflow/pull/73315/commits/72d150a27844a1831d36c46f4f2a3159d55a1582)



##########
task-sdk/src/airflow/sdk/execution_time/task_runner.py:
##########
@@ -1569,7 +1570,11 @@ def _on_term(signum, frame):
     try:
         # First, clear the xcom data sent from server
         if ti._ti_context_from_server and (keys_to_delete := 
ti._ti_context_from_server.xcom_keys_to_clear):
+            link_xcom_keys = {oe.xcom_key for oe in 
ti.task.operator_extra_links}
             for x in keys_to_delete:
+                if is_link_xcom_key(x, link_xcom_keys):
+                    # skip clearing this key as it is an operator link
+                    continue

Review Comment:
   [comements from 
ash](https://github.com/apache/airflow/pull/73315/commits/72d150a27844a1831d36c46f4f2a3159d55a1582)



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