rjgoyln commented on code in PR #73709:
URL: https://github.com/apache/airflow/pull/73709#discussion_r4120871665


##########
providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py:
##########
@@ -840,6 +840,35 @@ def _change_state(
             self.event_buffer[key] = state, None
             return
 
+        # Only pods this executor launched and is still tracking can be 
requeued; checking the
+        # in-memory attempt first avoids a metadata-db lookup for adopted or 
already-finalized pods.
+        attempt = self.pod_launch_attempts.get(key)
+        if attempt is not None and state == TaskInstanceState.FAILED and 
attempt.requeued_for_pod == pod_name:
+            # Kubernetes can emit several Failed events for one pod; we 
already requeued
+            # for this one, so ignore the duplicate before any pod API calls.
+            self.log.debug(
+                "Ignoring duplicate pre-execution failure for already-requeued 
pod %s/%s",
+                namespace,
+                pod_name,
+            )
+            return
+
+        # Completed pods adopted from a scheduler that is no longer alive are 
never tracked in
+        # self.running; their delete_pod call below is the intended cleanup 
(not a duplicate),
+        # and nothing is reported to the scheduler for them.
+        adopted_completed_pod = state == "completed"
+        if not adopted_completed_pod:
+            # The watcher can deliver the same completion event more than once 
for a
+            # TaskInstanceKey (e.g. mapped tasks). The first pass removes the 
key from
+            # self.running, so a repeat is a duplicate whose pod API calls 
already ran:
+            # drop it here, before issuing another delete_pod that would only 
404
+            # against the API server.
+            try:
+                self.running.remove(key)

Review Comment:
   `sync()` catches exceptions from `_change_state` and re-queues the result, 
while `delete_pod` re-raises non-404 `ApiException`s (such as 429 or 5xx 
responses, which seem relevant given the rate-limit handling elsewhere in this 
file). Because the removal from `self.running` now happens before the pod API 
calls, a retry after such a failure finds the key missing and treats the result 
as a duplicate. The pod is therefore never deleted, and the terminal state 
never reaches `event_buffer`, so the task instance receives no executor event.
   
   I verified this on the branch: when the first `delete_pod` call for a 
SUCCESS result raises `ApiException(status=429)`, `delete_pod` is called once 
and `event_buffer` remains empty. On `main`, it is called twice and the buffer 
contains `success`.
   
   The same issue occurs with `delete_worker_pods=False`: the retry can skip 
the `airflow_executor_done` patch that the watcher's 
`airflow_executor_done!=True` selector relies on.
   
   Would it make sense to wrap the two pod API calls and restore the key to 
`self.running` if either call raises?
   
   ```python
   except Exception:
       if not adopted_completed_pod:
           self.running.add(key)
       raise
   ```
   
   This keeps the duplicate-event deduplication while preserving the existing 
retry behavior.
   
   With this change, both new tests, both 
`test_sync_processes_completed_pods_once*` tests, and the rest of the executors 
suite remain green.
   



##########
providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py:
##########
@@ -840,6 +840,35 @@ def _change_state(
             self.event_buffer[key] = state, None
             return
 
+        # Only pods this executor launched and is still tracking can be 
requeued; checking the
+        # in-memory attempt first avoids a metadata-db lookup for adopted or 
already-finalized pods.
+        attempt = self.pod_launch_attempts.get(key)
+        if attempt is not None and state == TaskInstanceState.FAILED and 
attempt.requeued_for_pod == pod_name:
+            # Kubernetes can emit several Failed events for one pod; we 
already requeued
+            # for this one, so ignore the duplicate before any pod API calls.
+            self.log.debug(
+                "Ignoring duplicate pre-execution failure for already-requeued 
pod %s/%s",
+                namespace,
+                pod_name,
+            )
+            return
+
+        # Completed pods adopted from a scheduler that is no longer alive are 
never tracked in
+        # self.running; their delete_pod call below is the intended cleanup 
(not a duplicate),
+        # and nothing is reported to the scheduler for them.
+        adopted_completed_pod = state == "completed"

Review Comment:
   `"completed"` is now load-bearing for correctness in two separate places: 
`_adopt_completed_pods` sets it, and this branch uses it to decide cleanup and 
reporting. A typo on either side would therefore silently take an adopted pod 
down the wrong path instead of failing fast.
   
   `ADOPTED` in `kubernetes_executor_types.py` already follows this pattern. 
Would it make sense to add `COMPLETED = "completed"` alongside it and use the 
constant here as well?
   
   Non-blocking.
   



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