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]

Reply via email to