This is an automated email from the ASF dual-hosted git repository.

shahar1 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 dbac03d87dd Fix KubernetesExecutor not queuing tasks on Airflow 3.0.x 
(#73764)
dbac03d87dd is described below

commit dbac03d87dd98c93c71a815bc9a3635a81821681
Author: Anish Giri <[email protected]>
AuthorDate: Mon Sep 28 15:17:24 2026 -0500

    Fix KubernetesExecutor not queuing tasks on Airflow 3.0.x (#73764)
---
 .../kubernetes/executors/kubernetes_executor.py    | 12 +++++++-
 .../executors/test_kubernetes_executor.py          | 32 ++++++++++++++++++++++
 2 files changed, 43 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 dbb444d5a4c..713b6e83a26 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
@@ -59,7 +59,7 @@ from 
airflow.providers.cncf.kubernetes.kubernetes_helper_functions import (
     annotations_to_key,
 )
 from airflow.providers.cncf.kubernetes.pod_generator import PodGenerator
-from airflow.providers.cncf.kubernetes.version_compat import AIRFLOW_V_3_4_PLUS
+from airflow.providers.cncf.kubernetes.version_compat import 
AIRFLOW_V_3_1_PLUS, AIRFLOW_V_3_4_PLUS
 from airflow.providers.common.compat.sdk import Stats, conf
 from airflow.utils.helpers import prune_dict
 from airflow.utils.log.logging_mixin import remove_escape_codes
@@ -371,6 +371,16 @@ class KubernetesExecutor(BaseExecutor):
         self.pod_launch_attempts[key] = _PodLaunchAttempt(job=job)
         self.task_queue.put(job)
 
+    # TODO: Remove this once the minimum supported Airflow version is 3.1+ and 
defer to BaseExecutor.queue_workload.
+    if not AIRFLOW_V_3_1_PLUS:
+
+        def queue_workload(self, workload: workloads.All, session: Session | 
None) -> None:
+            from airflow.executors import workloads
+
+            if not isinstance(workload, workloads.ExecuteTask):
+                raise RuntimeError(f"{type(self)} cannot handle workloads of 
type {type(workload)}")
+            self.queued_tasks[workload.ti.key] = workload
+
     def _process_workloads(self, workloads: Sequence[workloads.All]) -> None:
         from airflow.executors.workloads import ExecuteTask
 
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 364eeecb9c9..4ac6245fa22 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
@@ -1740,6 +1740,38 @@ class TestKubernetesExecutor:
             key=key, command=[workload], queue="default", executor_config={}
         )
 
+    @pytest.mark.skipif(not AIRFLOW_V_3_0_PLUS, reason="Test requires Airflow 
3+")
+    def test_queue_workload_queues_execute_task(self):
+        """queue_workload must queue an ExecuteTask on every supported Airflow 
3 version.
+
+        On Airflow 3.0.x ``BaseExecutor.queue_workload`` raises 
unconditionally and the scheduler falls back
+        to ``queue_command`` for executors without their own override, which 
this executor cannot process.
+        """
+        from airflow.executors.workloads import ExecuteTask
+
+        executor = self.kubernetes_executor
+        key = TaskInstanceKey("dag", "task", "run_id", 1, -1)
+        workload = mock.Mock(spec=ExecuteTask)
+        workload.ti = mock.Mock()
+        workload.ti.key = key
+
+        if AIRFLOW_V_3_4_PLUS:
+            from airflow.executors.workloads.base import WorkloadType
+
+            workload.type = WorkloadType.EXECUTE_TASK
+            workload.key = key
+            task_queue = executor.executor_queues[WorkloadType.EXECUTE_TASK]
+        else:
+            task_queue = executor.queued_tasks
+
+        executor.queue_workload(workload, session=mock.MagicMock())
+
+        assert task_queue[key] is workload
+        if not AIRFLOW_V_3_1_PLUS:
+            from airflow.executors.base_executor import BaseExecutor
+
+            assert KubernetesExecutor.queue_workload is not 
BaseExecutor.queue_workload
+
     
@mock.patch("airflow.providers.cncf.kubernetes.executors.kubernetes_executor_utils.KubernetesJobWatcher")
     
@mock.patch("airflow.providers.cncf.kubernetes.kube_client.get_kube_client")
     def test_invalid_executor_config(self, mock_get_kube_client, 
mock_kubernetes_job_watcher):

Reply via email to