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

Reply via email to