ashb commented on code in PR #73916:
URL: https://github.com/apache/airflow/pull/73916#discussion_r4160069952
##########
airflow-core/src/airflow/executors/base_executor.py:
##########
@@ -641,11 +693,54 @@ def get_event_buffer(self, dag_ids=None) ->
dict[WorkloadKey, EventBufferValueTy
self.event_buffer = {}
else:
for key in list(self.event_buffer.keys()):
- if not isinstance(key, TaskInstanceKey) or key.dag_id in
dag_ids:
+ coordinates = self._task_coordinates.get(key) if
isinstance(key, TaskInstanceUuid) else None
+ if isinstance(key, TaskInstanceKey):
+ coordinates = key
+ if not isinstance(key, (TaskInstanceUuid, TaskInstanceKey)) or
(
+ coordinates is not None and coordinates.dag_id in dag_ids
+ ):
cleared_events[key] = self.event_buffer.pop(key)
+ task_queue = self.executor_queues.get(WorkloadType.EXECUTE_TASK, {})
+ for task_id, coordinates in list(self._task_coordinates.items()):
+ if not any(
+ key in container
+ for key in (task_id, coordinates)
+ for container in (self.running, task_queue, self.event_buffer)
+ ):
+ del self._task_coordinates[task_id]
+
return cleared_events
+ def _drain_events_with_task_ids(
+ self, dag_ids=None
+ ) -> tuple[dict[WorkloadKey, EventBufferValueType], dict[TaskInstanceUuid,
TaskInstanceKey]]:
+ """Drain events with captured attempt identities and coordinates."""
+ captured_coordinates = self._task_coordinates.copy()
+ task_ids: dict[TaskInstanceKey, TaskInstanceUuid | None] = {}
+ for task_id, coordinates in captured_coordinates.items():
+ task_ids[coordinates] = None if coordinates in task_ids else
task_id
Review Comment:
This requires deleting and recreating a run with identical task coordinates
while the executor still tracks the old attempt. An older Celery provider
doesn’t include enough identity in its events to distinguish those attempts
safely.
If the replacement reports its final state normally, that state wins. If its
worker dies before reporting, the existing task heartbeat timeout detects the
missing heartbeats and handles failure and retries. The consequence is delayed
recovery in a rare case.
Discarding the ambiguous event avoids potentially failing a healthy
replacement. The updated Celery provider in the second PR removes this
ambiguity by using UUIDs.
So in short "there is no perfect solution", this one "fails closed" in a
somewhat sensible pattern in the (unlikely?) case that this happens.
--
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]