From f60f841ac64b3696909d68db6d5831074d7df47c Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Guilherme=20Da=20Silva=20Gon=C3=A7alves?= Date: Sat, 11 Jul 2026 18:24:19 +0200 Subject: [PATCH 1/7] Prevent KubernetesExecutor from launching stale workloads --- .../executors/kubernetes_executor.py | 89 ++++++++++++++++++- .../executors/test_kubernetes_executor.py | 43 +++++++++ 2 files changed, 131 insertions(+), 1 deletion(-) diff --git a/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py index b75815df2fd91..83df21d5c9529 100644 --- a/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py +++ b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py @@ -398,6 +398,85 @@ def _process_workloads(self, workloads: Sequence[workloads.All]) -> None: self.execute_async(key=key, command=command, queue=queue, executor_config=executor_config) self.running.add(key) + @provide_session + def _should_create_pod_for_job( + self, + task: KubernetesJob, + *, + session: Session = NEW_SESSION, + ) -> bool: + """ + Check whether an executor job still represents the current queued task instance. + + The scheduler creates an ``ExecuteTask`` workload while the task instance is queued, but the + Kubernetes pod may be created much later, for example after API-server throttling or quota + failures. In an HA scheduler deployment, the task instance may have been retried, cleared, or + otherwise replaced before this executor gets another chance to create the pod. Revalidating the + immutable task instance id and launch ownership here prevents an obsolete workload from creating + a stale worker pod. + """ + from airflow.executors.workloads import ExecuteTask + from airflow.models.taskinstance import TaskInstance + + if not task.command or not isinstance(task.command[0], ExecuteTask): + return True + + workload = task.command[0] + workload_ti = workload.ti + try: + scheduler_job_id = int(self.scheduler_job_id) if self.scheduler_job_id is not None else None + except ValueError: + self.log.debug( + "Skipping stale Kubernetes workload check because scheduler_job_id %r is not numeric", + self.scheduler_job_id, + ) + return True + + ti = session.execute( + select( + TaskInstance.id, + TaskInstance.state, + TaskInstance.try_number, + TaskInstance.queued_by_job_id, + ).where(TaskInstance.id == workload_ti.id) + ).one_or_none() + if ti is None: + self.log.info( + "Dropping stale Kubernetes workload for %s because task instance id %s no longer exists", + task.key, + workload_ti.id, + ) + return False + + _, state, try_number, queued_by_job_id = ti + if ( + state == TaskInstanceState.QUEUED + and try_number == workload_ti.try_number + and queued_by_job_id == scheduler_job_id + ): + return True + + self.log.info( + "Dropping stale Kubernetes workload for %s because current task instance state does not " + "match the queued workload. task_instance_id=%s, state=%s, try_number=%s, " + "queued_by_job_id=%s, workload_try_number=%s, scheduler_job_id=%s", + task.key, + workload_ti.id, + state, + try_number, + queued_by_job_id, + workload_ti.try_number, + scheduler_job_id, + ) + return False + + def _discard_stale_pod_creation_task(self, task: KubernetesJob) -> None: + """Remove executor bookkeeping for a stale job that will not create a pod.""" + self.running.discard(task.key) + if self.event_buffer.get(task.key) == (TaskInstanceState.QUEUED, self.scheduler_job_id): + self.event_buffer.pop(task.key, None) + Stats.incr("kubernetes_executor.stale_workload_dropped") + def sync(self) -> None: """Synchronize task state.""" if TYPE_CHECKING: @@ -486,6 +565,9 @@ def _create_pods_sequentially(self) -> None: task: KubernetesJob = self.task_queue.get_nowait() created += 1 try: + if not self._should_create_pod_for_job(task): + self._discard_stale_pod_creation_task(task) + continue self.kube_scheduler.run_next(task) self.task_publish_retries.pop(task.key, None) except ( @@ -517,7 +599,12 @@ def _create_pods_concurrently(self) -> None: jobs: list[KubernetesJob] = [] with contextlib.suppress(Empty): for _ in range(self.kube_config.worker_pods_creation_batch_size): - jobs.append(self.task_queue.get_nowait()) + task = self.task_queue.get_nowait() + if not self._should_create_pod_for_job(task): + self._discard_stale_pod_creation_task(task) + self.task_queue.task_done() + continue + jobs.append(task) if not jobs: return start: float = time.monotonic() diff --git a/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/executors/test_kubernetes_executor.py b/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/executors/test_kubernetes_executor.py index 5f1c03e7d4524..3d39864634737 100644 --- a/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/executors/test_kubernetes_executor.py +++ b/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/executors/test_kubernetes_executor.py @@ -1088,6 +1088,49 @@ def test_skip_pod_creation_on_create_pods_after( finally: kubernetes_executor.end() + @pytest.mark.db_test + @pytest.mark.skipif(not AIRFLOW_V_3_0_PLUS, reason="workloads are used on Airflow 3+") + @mock.patch("airflow.providers.cncf.kubernetes.executors.kubernetes_executor_utils.KubernetesJobWatcher") + @mock.patch("airflow.providers.cncf.kubernetes.kube_client.get_kube_client") + def test_sync_drops_stale_execute_task_workload_before_pod_creation( + self, + mock_get_kube_client, + mock_kubernetes_job_watcher, + create_task_instance, + session, + ): + """A delayed Kubernetes workload should not create a pod after the DB task moved on.""" + from airflow.executors.workloads import ExecuteTask + + executor = self.kubernetes_executor + executor.start() + try: + ti = create_task_instance(state=TaskInstanceState.QUEUED) + ti.queued_by_job_id = executor.job_id + session.merge(ti) + session.commit() + + workload = ExecuteTask.make(ti) + executor.queue_workload(workload, session=session) + executor._process_workloads([workload]) + + ti.state = TaskInstanceState.SUCCESS + session.merge(ti) + session.commit() + + assert executor.kube_scheduler is not None + executor.kube_scheduler.run_next = mock.Mock() + + executor.sync() + + executor.kube_scheduler.run_next.assert_not_called() + assert executor.task_queue is not None + assert executor.task_queue.empty() + assert ti.key not in executor.running + assert ti.key not in executor.event_buffer + finally: + executor.end() + @pytest.mark.skipif( AirflowKubernetesScheduler is None, reason="kubernetes python package is not installed" ) From 9fe368bea3d3dca979d5ffe6c8a3d9502493370b Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Guilherme=20Da=20Silva=20Gon=C3=A7alves?= Date: Mon, 20 Jul 2026 15:14:45 +0200 Subject: [PATCH 2/7] Fix stale workload CI failures --- .../executors/kubernetes_executor.py | 29 ++++++++++++------- .../metrics/metrics_template.yaml | 6 ++++ 2 files changed, 25 insertions(+), 10 deletions(-) diff --git a/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py index 83df21d5c9529..44e685f92f28f 100644 --- a/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py +++ b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py @@ -398,13 +398,7 @@ def _process_workloads(self, workloads: Sequence[workloads.All]) -> None: self.execute_async(key=key, command=command, queue=queue, executor_config=executor_config) self.running.add(key) - @provide_session - def _should_create_pod_for_job( - self, - task: KubernetesJob, - *, - session: Session = NEW_SESSION, - ) -> bool: + def _should_create_pod_for_job(self, task: KubernetesJob) -> bool: """ Check whether an executor job still represents the current queued task instance. @@ -415,13 +409,28 @@ def _should_create_pod_for_job( immutable task instance id and launch ownership here prevents an obsolete workload from creating a stale worker pod. """ - from airflow.executors.workloads import ExecuteTask - from airflow.models.taskinstance import TaskInstance + try: + from airflow.executors.workloads import ExecuteTask + except ImportError: + # Compatibility with older Airflow versions tested by provider compatibility jobs. + return True if not task.command or not isinstance(task.command[0], ExecuteTask): return True - workload = task.command[0] + return self._should_create_pod_for_execute_task(task, task.command[0]) + + @provide_session + def _should_create_pod_for_execute_task( + self, + task: KubernetesJob, + workload: Any, + *, + session: Session = NEW_SESSION, + ) -> bool: + """Check that an ``ExecuteTask`` workload still owns the queued task instance row.""" + from airflow.models.taskinstance import TaskInstance + workload_ti = workload.ti try: scheduler_job_id = int(self.scheduler_job_id) if self.scheduler_job_id is not None else None diff --git a/shared/observability/src/airflow_shared/observability/metrics/metrics_template.yaml b/shared/observability/src/airflow_shared/observability/metrics/metrics_template.yaml index 536e439f876e2..665ab7aababf2 100644 --- a/shared/observability/src/airflow_shared/observability/metrics/metrics_template.yaml +++ b/shared/observability/src/airflow_shared/observability/metrics/metrics_template.yaml @@ -619,6 +619,12 @@ metrics: legacy_name: "-" name_variables: ["status"] + - name: "kubernetes_executor.stale_workload_dropped" + description: "Number of stale KubernetesExecutor workloads dropped before worker pod creation." + type: "counter" + legacy_name: "-" + name_variables: [] + - name: "kubernetes_executor.pod_deletion_status" description: "Number of Kubernetes delete_namespaced_pod calls from the Kubernetes Executor, tagged by HTTP response status (``200`` on success, the ``ApiException`` status code on failure)." From a0165e991e89059262c1d7d84a4c521d2cddbc09 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Guilherme=20Da=20Silva=20Gon=C3=A7alves?= Date: Tue, 21 Jul 2026 14:10:46 +0200 Subject: [PATCH 3/7] Fix stale workload lookup for UUID task ids --- .../providers/cncf/kubernetes/executors/kubernetes_executor.py | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py index 44e685f92f28f..71be95703ce77 100644 --- a/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py +++ b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py @@ -441,13 +441,14 @@ def _should_create_pod_for_execute_task( ) return True + task_instance_id = str(workload_ti.id) ti = session.execute( select( TaskInstance.id, TaskInstance.state, TaskInstance.try_number, TaskInstance.queued_by_job_id, - ).where(TaskInstance.id == workload_ti.id) + ).where(TaskInstance.id == task_instance_id) ).one_or_none() if ti is None: self.log.info( From 8f2ee1ab47dca6fb21bbb48a19bb565b9834c91c Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Guilherme=20Da=20Silva=20Gon=C3=A7alves?= Date: Tue, 21 Jul 2026 18:43:12 +0200 Subject: [PATCH 4/7] Cast task instance ids for stale workload lookup --- .../cncf/kubernetes/executors/kubernetes_executor.py | 5 ++--- 1 file changed, 2 insertions(+), 3 deletions(-) diff --git a/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py index 71be95703ce77..b6d0f5e155b5d 100644 --- a/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py +++ b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py @@ -41,7 +41,7 @@ from deprecated import deprecated from kubernetes.dynamic import DynamicClient -from sqlalchemy import select +from sqlalchemy import String, cast, select from airflow.exceptions import AirflowProviderDeprecationWarning from airflow.executors.base_executor import BaseExecutor @@ -441,14 +441,13 @@ def _should_create_pod_for_execute_task( ) return True - task_instance_id = str(workload_ti.id) ti = session.execute( select( TaskInstance.id, TaskInstance.state, TaskInstance.try_number, TaskInstance.queued_by_job_id, - ).where(TaskInstance.id == task_instance_id) + ).where(cast(TaskInstance.id, String) == str(workload_ti.id)) ).one_or_none() if ti is None: self.log.info( From 711a520e3525005c4e0a1df1da86ed1c5e0d3c66 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Guilherme=20Da=20Silva=20Gon=C3=A7alves?= Date: Wed, 22 Jul 2026 11:15:54 +0200 Subject: [PATCH 5/7] Remove stale workload metric from provider PR --- .../cncf/kubernetes/executors/kubernetes_executor.py | 1 - .../observability/metrics/metrics_template.yaml | 6 ------ 2 files changed, 7 deletions(-) diff --git a/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py index b6d0f5e155b5d..f274cce5eb590 100644 --- a/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py +++ b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py @@ -484,7 +484,6 @@ def _discard_stale_pod_creation_task(self, task: KubernetesJob) -> None: self.running.discard(task.key) if self.event_buffer.get(task.key) == (TaskInstanceState.QUEUED, self.scheduler_job_id): self.event_buffer.pop(task.key, None) - Stats.incr("kubernetes_executor.stale_workload_dropped") def sync(self) -> None: """Synchronize task state.""" diff --git a/shared/observability/src/airflow_shared/observability/metrics/metrics_template.yaml b/shared/observability/src/airflow_shared/observability/metrics/metrics_template.yaml index 665ab7aababf2..536e439f876e2 100644 --- a/shared/observability/src/airflow_shared/observability/metrics/metrics_template.yaml +++ b/shared/observability/src/airflow_shared/observability/metrics/metrics_template.yaml @@ -619,12 +619,6 @@ metrics: legacy_name: "-" name_variables: ["status"] - - name: "kubernetes_executor.stale_workload_dropped" - description: "Number of stale KubernetesExecutor workloads dropped before worker pod creation." - type: "counter" - legacy_name: "-" - name_variables: [] - - name: "kubernetes_executor.pod_deletion_status" description: "Number of Kubernetes delete_namespaced_pod calls from the Kubernetes Executor, tagged by HTTP response status (``200`` on success, the ``ApiException`` status code on failure)." From b99c29d40c906147607bda4cfd696d43529cac36 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Guilherme=20Da=20Silva=20Gon=C3=A7alves?= Date: Thu, 30 Jul 2026 20:27:41 +0200 Subject: [PATCH 6/7] Keep stale workload lookup indexed --- .../cncf/kubernetes/executors/kubernetes_executor.py | 8 +++++--- 1 file changed, 5 insertions(+), 3 deletions(-) diff --git a/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py index f274cce5eb590..65ce5878e9177 100644 --- a/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py +++ b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py @@ -38,10 +38,11 @@ from itertools import chain from queue import Empty, Queue from typing import TYPE_CHECKING, Any +from uuid import UUID from deprecated import deprecated from kubernetes.dynamic import DynamicClient -from sqlalchemy import String, cast, select +from sqlalchemy import select from airflow.exceptions import AirflowProviderDeprecationWarning from airflow.executors.base_executor import BaseExecutor @@ -77,6 +78,7 @@ from airflow._shared.logging.remote import RawLogStream, StreamingLogResponse from airflow.cli.cli_config import GroupCommand from airflow.executors import workloads + from airflow.executors.workloads import ExecuteTask from airflow.models.taskinstance import TaskInstance from airflow.models.taskinstancekey import TaskInstanceKey from airflow.providers.cncf.kubernetes.executors.kubernetes_executor_utils import ( @@ -424,7 +426,7 @@ def _should_create_pod_for_job(self, task: KubernetesJob) -> bool: def _should_create_pod_for_execute_task( self, task: KubernetesJob, - workload: Any, + workload: ExecuteTask, *, session: Session = NEW_SESSION, ) -> bool: @@ -447,7 +449,7 @@ def _should_create_pod_for_execute_task( TaskInstance.state, TaskInstance.try_number, TaskInstance.queued_by_job_id, - ).where(cast(TaskInstance.id, String) == str(workload_ti.id)) + ).where(TaskInstance.id == UUID(str(workload_ti.id))) ).one_or_none() if ti is None: self.log.info( From 50d28e56bffd099a242da074bc806385ab9c948b Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Guilherme=20Da=20Silva=20Gon=C3=A7alves?= Date: Fri, 31 Jul 2026 00:00:22 +0200 Subject: [PATCH 7/7] Keep stale workload lookup compatible across Airflow versions --- .../kubernetes/executors/kubernetes_executor.py | 15 +++++++++++++-- 1 file changed, 13 insertions(+), 2 deletions(-) diff --git a/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py index 65ce5878e9177..bf969c25fd27c 100644 --- a/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py +++ b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py @@ -38,7 +38,6 @@ from itertools import chain from queue import Empty, Queue from typing import TYPE_CHECKING, Any -from uuid import UUID from deprecated import deprecated from kubernetes.dynamic import DynamicClient @@ -443,13 +442,25 @@ def _should_create_pod_for_execute_task( ) return True + # Bind the id as the mapped column's own Python type so the predicate stays sargable and + # uses the primary-key index. The column is a native ``Uuid`` on Airflow 3.2+ but a ``String`` + # on older versions exercised by provider compatibility jobs, and SQLite rejects binding a + # ``UUID`` object against a string column. + try: + ti_id_python_type = TaskInstance.id.type.python_type + except NotImplementedError: + ti_id_python_type = str + ti_id = ( + ti_id_python_type(str(workload_ti.id)) if ti_id_python_type is not str else str(workload_ti.id) + ) + ti = session.execute( select( TaskInstance.id, TaskInstance.state, TaskInstance.try_number, TaskInstance.queued_by_job_id, - ).where(TaskInstance.id == UUID(str(workload_ti.id))) + ).where(TaskInstance.id == ti_id) ).one_or_none() if ti is None: self.log.info(