This is an automated email from the ASF dual-hosted git repository.
potiuk pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/airflow.git
The following commit(s) were added to refs/heads/main by this push:
new 90972c41dfe Prevent KubernetesExecutor from launching stale workloads
(#69762)
90972c41dfe is described below
commit 90972c41dfe69720a74a1b376d710abf626169c3
Author: Guilherme Da Silva Gonçalves <[email protected]>
AuthorDate: Wed Sep 23 02:35:51 2026 +0200
Prevent KubernetesExecutor from launching stale workloads (#69762)
* Prevent KubernetesExecutor from launching stale workloads
* Fix stale workload CI failures
* Fix stale workload lookup for UUID task ids
* Cast task instance ids for stale workload lookup
* Remove stale workload metric from provider PR
* Keep stale workload lookup indexed
* Keep stale workload lookup compatible across Airflow versions
* Simplify redundant ti_id coercion ternary in KubernetesExecutor
* Make stale-workload test independent of queue_workload on Airflow 3.0
BaseExecutor.queue_workload only accepts ExecuteTask from Airflow 3.1, so
the provider compatibility job on Airflow 3.0.6 failed in the test's setup
rather than in the behaviour under test. Enqueueing the pod-creation job
directly keeps the Airflow 3.0 coverage.
Generated-by: Claude Opus 5
---------
Co-authored-by: Jarek Potiuk <[email protected]>
---
.../kubernetes/executors/kubernetes_executor.py | 108 ++++++++++++++++++++-
.../executors/test_kubernetes_executor.py | 45 +++++++++
2 files changed, 152 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 3216ed232b8..dbb444d5a4c 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
@@ -80,6 +80,7 @@ if TYPE_CHECKING:
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 (
@@ -391,6 +392,103 @@ class KubernetesExecutor(BaseExecutor):
self.execute_async(key=key, command=command, queue=queue,
executor_config=executor_config)
self.running.add(key)
+ def _should_create_pod_for_job(self, task: KubernetesJob) -> 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.
+ """
+ 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
+
+ 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: ExecuteTask,
+ *,
+ 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
+ except ValueError:
+ self.log.debug(
+ "Skipping stale Kubernetes workload check because
scheduler_job_id %r is not numeric",
+ self.scheduler_job_id,
+ )
+ 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))
+
+ ti = session.execute(
+ select(
+ TaskInstance.id,
+ TaskInstance.state,
+ TaskInstance.try_number,
+ TaskInstance.queued_by_job_id,
+ ).where(TaskInstance.id == 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)
+
def sync(self) -> None:
"""Synchronize task state."""
if TYPE_CHECKING:
@@ -482,6 +580,9 @@ class KubernetesExecutor(BaseExecutor):
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 (
@@ -513,7 +614,12 @@ class KubernetesExecutor(BaseExecutor):
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 15b95ab4d1c..364eeecb9c9 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
@@ -1083,6 +1083,51 @@ class TestKubernetesExecutor:
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)
+ # Enqueue the pod-creation job directly:
`BaseExecutor.queue_workload` only accepts
+ # `ExecuteTask` from Airflow 3.1, and the provider compat jobs
also run this on 3.0.
+ executor.execute_async(key=ti.key, command=[workload],
queue=ti.queue, executor_config={})
+ executor.running.add(ti.key)
+
+ 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"
)