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:
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:
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:
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:
-
Define a DAG with a single KubernetesPodOperator running a long container:
KubernetesPodOperator(
task_id="sleeper",
image="alpine:3",
cmds=["sh", "-c", "sleep 600 && echo done"],
on_finish_action="delete_pod", # the default
get_logs=True,
)
-
Trigger the DAG. Wait until the pod reaches Running and the operator is streaming its logs.
-
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.
-
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.:
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
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:
Are you willing to submit PR?
Code of Conduct
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
KubernetesPodOperatorreceivesSIGTERMwhileexecute()is in flight, the task instance is committed assuccesseven though the operator's pod was destroyed mid-run and its work never completed.The two interacting pieces:
(a)
SIGTERMinvokeson_kill()but does not fail the task.airflow/sdk/execution_time/task_runner.py:1544: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, insideKubernetesPodOperator.execute(). Theexcept AirflowTaskTerminatedbranch attask_runner.py:1655carries the comment "these exceptions should ideally never be thrown", and on this path nothing throws it.(b)
KubernetesPodOperator.on_kill()destroys the work, thencleanup()silently skips the failure raise.providers/cncf/kubernetes/operators/pod.py:1556:providers/cncf/kubernetes/operators/pod.py:1339:cleanup()is the only thing that raisesAirflowExceptionfor a pod that did not reachSucceeded. The_killedguard (introduced in d3b4a91, "fix: Avoid retrying after KubernetesPodOperator has been marked as failed (#36749)") skips it.That guard's premise —
_killedmeans the TI is already being marked failed elsewhere — holds when a user marks a task failed in the UI. It does not hold when theSIGTERMcomes from a pod eviction or any other external termination, because in that case nothing else marks the TI failed.Net result:
execute_sync()returnsNonenormally, no exception ever reaches_run_task_and_map_outcome, and the task is committedsuccesswith a normalend_date.Steps to reproduce
Minimal form — does not require an eviction, any executor that runs the task in its own process will do:
Define a DAG with a single
KubernetesPodOperatorrunning a long container:Trigger the DAG. Wait until the pod reaches
Runningand the operator is streaming its logs.Send
SIGTERMto the task process (not the pod):Equivalently, under
KubernetesExecutor:kubectl delete pod <the KE worker pod>, or cordon/drain the node it is on — anything that delivers a gracefulSIGTERMto the worker.Observe:
alpinepod is deleted immediately byon_kill(), having never printeddone;success.How it shows up in production
In our deployment this fires naturally roughly 0.73×/day (33 occurrences over 45 days of
task_instanceretention). The trigger is an EKS managed-node-group rolling drain:eks:node-managerissuescreate pods/evictionagainst the KubernetesExecutor worker pod. From the EKS control-plane audit log: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:
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:Pod %s has phase %s—PodManager.await_pod_completionlogs this on every non-terminal poll, so its absence shows the loop exited on the first read;Deleting pod: %sand noSkipping deleting pod: %s— provingcleanup()returned at the_killedguard before reachingprocess_pod_deletion;A healthy run of the same task always logs both
Pod ... has phase RunningandDeleting pod: ...beforePost 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 configuredretries), notsuccess. 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
_killedguard incleanup(). Only skip the failure raise when the pod actually succeeded, e.g.: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_termactually terminate the task. Afterti.task.on_kill(),_on_termcould raiseAirflowTaskTerminated(or set a flag checked immediately afterexecute()returns) so the existingexcept AirflowTaskTerminatedhandler attask_runner.py:1655fires and the TI is markedFAILED. 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 whoseon_kill()tears down out-of-process work, not justKubernetesPodOperator.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
Relevant only as context for how the
SIGTERMarrives — the bug reproduces regardless of these settings. Note thatsafe-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:
Are you willing to submit PR?
Code of Conduct