ashb commented on code in PR #73315:
URL: https://github.com/apache/airflow/pull/73315#discussion_r4103188888
##########
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:
Isn't there a `.get_one()`?
##########
task-sdk/tests/task_sdk/execution_time/test_task_runner.py:
##########
@@ -2901,14 +2901,16 @@ def execute(self, context):
state=TaskInstanceState.SUCCESS,
context=runtime_ti.get_template_context(),
)
- mock_xcom_set.assert_called_once_with(
- key="_link_AirflowLink",
- value="https://airflow.apache.org",
- dag_id=runtime_ti.dag_id,
- task_id=runtime_ti.task_id,
- run_id=runtime_ti.run_id,
- map_index=runtime_ti.map_index,
- )
+ assert mock_xcom_set.mock_calls == [
+ call(
+ key="_link_AirflowLink__try_1",
Review Comment:
Does the storage tab in the UI already filter out the se
`^_link_AirflowLink` items?
##########
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:
This should be done on the server side, not the client side I feel (i.e.
server shouldn't tell client to clear these)
--
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]