shahar1 commented on code in PR #73916:
URL: https://github.com/apache/airflow/pull/73916#discussion_r4157602910


##########
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:
   With `CeleryExecutor`, this can discard every event from the live attempt. 
If you mark a run as failed, delete it, and create a new run with the same 
run_id, both attempts have the same task coordinates. The PR treats the key as 
ambiguous and drops the new attempt’s `QUEUED` and `FAILED` events. The task 
can then remain `RUNNING` until the heartbeat timeout.
   Reproduced this with a Celery worker: after killing the new attempt’s 
worker, the PR left the task running for about 300 seconds; on main, the 
failure event marked it failed within about a second. The old attempt can’t 
emit an event in this case because Celery replaces its stored result when the 
new attempt is registered.
   Could the bridge use the newest registration for a reused key instead of 
dropping both events?



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