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