adrianrego opened a new issue, #73006:
URL: https://github.com/apache/airflow/issues/73006

   ### Under which category would you file this issue?
   
   Providers
   
   ### Apache Airflow version
   
   3.3.1
   
   ### What happened and how to reproduce it?
   
   ### Issue Description
   
   When the task process running a `KubernetesPodOperator` receives `SIGTERM` 
while `execute()` is in flight, the task instance is committed as **`success`** 
even though the operator's pod was destroyed mid-run and its work never 
completed.
   
   The two interacting pieces:
   
   **(a) `SIGTERM` invokes `on_kill()` but does not fail the task.**
   `airflow/sdk/execution_time/task_runner.py:1544`:
   
   ```python
   def _on_term(signum, frame):
       pid = os.getpid()
       if pid != parent_pid:
           return
   
       ti.task.on_kill()
   
   signal.signal(signal.SIGTERM, _on_term)
   ```
   
   The handler calls `on_kill()` and returns. It does not raise, does not set a 
terminating flag, and does not exit — so execution resumes exactly where it was 
interrupted, inside `KubernetesPodOperator.execute()`. The `except 
AirflowTaskTerminated` branch at `task_runner.py:1655` carries the comment 
*"these exceptions should ideally never be thrown"*, and on this path nothing 
throws it.
   
   **(b) `KubernetesPodOperator.on_kill()` destroys the work, then `cleanup()` 
silently skips the failure raise.**
   
   `providers/cncf/kubernetes/operators/pod.py:1556`:
   
   ```python
   def on_kill(self) -> None:
       self._killed = True
       ...
       self.client.delete_namespaced_pod(**kwargs)   # deletes the pod doing 
the actual work
   ```
   
   `providers/cncf/kubernetes/operators/pod.py:1339`:
   
   ```python
   def cleanup(self, pod, remote_pod, xcom_result=None, context=None) -> None:
       # Skip cleaning the pod in the following scenarios.
       # 1. If a task got marked as failed, "on_kill" method would be called ...
       # 2. remote pod is null (ex: pod creation failed)
       if self._killed or not remote_pod:
           return
   ```
   
   `cleanup()` is the only thing that raises `AirflowException` for a pod that 
did not reach `Succeeded`. The `_killed` guard (introduced in d3b4a91, *"fix: 
Avoid retrying after KubernetesPodOperator has been marked as failed 
(#36749)"*) skips it.
   
   That guard's premise — `_killed` means the TI is already being marked failed 
elsewhere — holds when a user marks a task failed in the UI. It does **not** 
hold when the `SIGTERM` comes from a pod eviction or any other external 
termination, because in that case nothing else marks the TI failed.
   
   **Net result:** `execute_sync()` returns `None` normally, no exception ever 
reaches `_run_task_and_map_outcome`, and the task is committed `success` with a 
normal `end_date`.
   
   ### Steps to reproduce
   
   Minimal form — does not require an eviction, any executor that runs the task 
in its own process will do:
   
   1. Define a DAG with a single `KubernetesPodOperator` running a long 
container:
   
      ```python
      KubernetesPodOperator(
          task_id="sleeper",
          image="alpine:3",
          cmds=["sh", "-c", "sleep 600 && echo done"],
          on_finish_action="delete_pod",   # the default
          get_logs=True,
      )
      ```
   
   2. Trigger the DAG. Wait until the pod reaches `Running` and the operator is 
streaming its logs.
   3. Send `SIGTERM` to the **task process** (not the pod):
   
      ```
      kill -TERM <pid of the airflow task-runner process>
      ```
   
      Equivalently, under `KubernetesExecutor`: `kubectl delete pod <the KE 
worker pod>`, or cordon/drain the node it is on — anything that delivers a 
graceful `SIGTERM` to the worker.
   
   4. Observe:
      - the `alpine` pod is deleted immediately by `on_kill()`, having never 
printed `done`;
      - the task instance is committed **`success`**.
   
   ### How it shows up in production
   
   In our deployment this fires naturally roughly **0.73×/day** (33 occurrences 
over 45 days of `task_instance` retention). The trigger is an EKS 
managed-node-group rolling drain: `eks:node-manager` issues `create 
pods/eviction` against the KubernetesExecutor worker pod. From the EKS 
control-plane audit log:
   
   ```
   23:49:43.038  create  pods/eviction  <ke-worker-pod>       
user=eks:node-manager
   23:49:43.048  delete  pods           <kpo-child-pod>       
user=<airflow-worker SA>   # on_kill(), 10 ms later
   23:49:51.052  delete  pods           <ke-worker-pod>       
user=system:node:<node>
   23:49:52.071  patch   pods           <ke-worker-pod>       user=<scheduler 
SA>  rc=404  (x3)
   ```
   
   The task log fingerprint is distinctive and is the easiest way for others to 
recognize this. Every non-container line from one affected attempt, in full:
   
   ```
   23:50:35.863  Pod has reached Running phase before launch timeout
   23:54:52.093  ::group::Post Execute
   23:54:52.112  ::endgroup::
   ```
   
   Four minutes of container log streaming, then straight to `Post Execute` — 
approximately 20–30 ms after the final container log line. Note what is 
**missing**:
   
   - no `Pod %s has phase %s` — `PodManager.await_pod_completion` logs this on 
every non-terminal poll, so its absence shows the loop exited on the first read;
   - no `Deleting pod: %s` and no `Skipping deleting pod: %s` — proving 
`cleanup()` returned at the `_killed` guard before reaching 
`process_pod_deletion`;
   - no exception, no traceback, no warning.
   
   A healthy run of the same task always logs both `Pod ... has phase Running` 
and `Deleting pod: ...` before `Post Execute`.
   
   The work loss is real and varies with where the kill lands: in one case the 
container had written 6 of its 8 output files and was mid-upload of the 7th; in 
others it had written none. Duration of the false-success attempt ranged from 
11% to 228% of the same task's eventual successful runtime, consistent with a 
kill at a uniformly random point.
   
   ---
   
   
   ### What you think should happen instead?
   
   The task should be marked **`failed`** (and therefore be eligible for its 
configured `retries`), not `success`. Silently reporting success for a task 
whose pod was destroyed mid-run means downstream tasks consume incomplete 
output with no signal at all — the failure is undetectable unless the DAG 
happens to independently verify the operator's side effects.
   
   Two candidate fixes; the second is broader and probably the more correct one:
   
   **(a) Provider — narrow the `_killed` guard in `cleanup()`.** Only skip the 
failure raise when the pod actually succeeded, e.g.:
   
   ```python
   if not remote_pod:
       return
   if self._killed and remote_pod.status.phase == PodPhase.SUCCEEDED:
       return
   ```
   
   This preserves the intent of #36749 (don't re-cleanup / don't trigger a 
retry for a task the user deliberately marked failed) while restoring the 
failure signal when the pod did not succeed. It would need care so that a 
UI-initiated "mark failed" does not start producing unwanted retries — that is 
exactly the regression #36749 fixed.
   
   **(b) Task SDK — make `_on_term` actually terminate the task.** After 
`ti.task.on_kill()`, `_on_term` could raise `AirflowTaskTerminated` (or set a 
flag checked immediately after `execute()` returns) so the existing `except 
AirflowTaskTerminated` handler at `task_runner.py:1655` fires and the TI is 
marked `FAILED`. That handler's own comment says such exceptions "should 
ideally never be thrown" — but on the SIGTERM path there is currently no 
mechanism that throws one, so a signal-terminated task can complete as success 
for *any* operator whose `on_kill()` tears down out-of-process work, not just 
`KubernetesPodOperator`.
   
   Either fix alone resolves the KPO case. (b) additionally covers other 
operators with the same shape.
   
   
   ### Operating System
   
   Ubuntu 24.04.5 LTS (container image, Python 3.12.14)
   
   ### Deployment
   
   Official Apache Airflow Helm Chart
   
   ### Apache Airflow Provider(s)
   
   amazon, celery, cncf-kubernetes, google, standard, pagerduty, http
   
   ### Versions of Apache Airflow Providers
   
   apache-airflow==3.3.1
   apache-airflow-core==3.3.1
   apache-airflow-task-sdk==1.3.1
   apache-airflow-providers-cncf-kubernetes==10.21.1
   apache-airflow-providers-celery==3.23.1
   apache-airflow-providers-standard==1.18.0
   
   
   ### Official Helm Chart version
   
   1.22.0 (latest released)
   
   ### Kubernetes Version
   
   1.34
   
   ### Helm Chart configuration
   
   ```yaml
   executor: "KubernetesExecutor,CeleryExecutor"
   
   config:
     kubernetes_executor:
       delete_worker_pods: "False"
       delete_worker_pods_on_failure: "False"
   
   workers:
     labels:
       executor: kubernetes
     podAnnotations:
       cluster-autoscaler.kubernetes.io/safe-to-evict: "false"
   ```
   Relevant only as context for how the `SIGTERM` arrives — the bug reproduces 
regardless of these settings. Note that `safe-to-evict: "false"` is honored by 
cluster-autoscaler only; an Eviction API call from the cloud provider's node 
manager is not affected by it, and there is no PodDisruptionBudget on 
KubernetesExecutor worker pods.
   
   
   
   ### Docker Image customizations
   
   _No response_
   
   ### Anything else?
   
   **Possibly related, but distinct:**
   - #53015 — retry/reattach behavior *after* a pod eviction (evicted pod 
reused via label selector). Different failure; that one at least fails loudly.
   - #21420 / discussion #21558 — tasks marked SUCCESS then immediately FAILED 
under memory pressure.
   - discussion #45830 — a UI-marked-failed KPO task turning green when its pod 
later completes.
   
   
   ### Are you willing to submit PR?
   
   - [ ] Yes I am willing to submit a PR!
   
   ### Code of Conduct
   
   - [x] I agree to follow this project's [Code of 
Conduct](https://github.com/apache/airflow/blob/main/CODE_OF_CONDUCT.md)
   


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