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]