diff --git a/docs/spelling_wordlist.txt b/docs/spelling_wordlist.txt index f39e41ee52876..397b6ccd9c806 100644 --- a/docs/spelling_wordlist.txt +++ b/docs/spelling_wordlist.txt @@ -959,6 +959,7 @@ kinesis kinit kms knownHosts +kpo krb Kube kube diff --git a/providers/cncf/kubernetes/docs/changelog.rst b/providers/cncf/kubernetes/docs/changelog.rst index 63659c31690c4..add40f7c34344 100644 --- a/providers/cncf/kubernetes/docs/changelog.rst +++ b/providers/cncf/kubernetes/docs/changelog.rst @@ -45,6 +45,7 @@ Changelog Features ~~~~~~~~ +* ``Add optional KubernetesPodOperator zombie pod cleanup`` * ``Decide pod_template and image based on Coordinator for lang-SDK tasks on KubernetesExecutor (#68713)`` * ``Add team_name tags to Kubernetes executor metrics (#69046)`` diff --git a/providers/cncf/kubernetes/provider.yaml b/providers/cncf/kubernetes/provider.yaml index 9e21164779089..7a369ea74e118 100644 --- a/providers/cncf/kubernetes/provider.yaml +++ b/providers/cncf/kubernetes/provider.yaml @@ -342,6 +342,35 @@ config: type: string example: ~ default: "False" + kpo_zombie_pod_cleanup_enabled: + description: | + If True, periodically delete KubernetesPodOperator pods that no longer have an active matching + task instance. + version_added: ~ + type: boolean + example: ~ + default: "False" + kpo_zombie_pod_cleanup_interval: + description: | + How often, in seconds, KubernetesPodOperator zombie pod cleanup runs when enabled. + version_added: ~ + type: integer + example: ~ + default: "300" + kpo_zombie_pod_cleanup_max_deletes_per_loop: + description: | + Maximum number of KubernetesPodOperator zombie pods to delete during one cleanup loop. + version_added: ~ + type: integer + example: ~ + default: "100" + kpo_zombie_pod_deletion_grace_period_seconds: + description: | + Grace period, in seconds, used when deleting KubernetesPodOperator zombie pods. + version_added: ~ + type: integer + example: ~ + default: "5" worker_pod_pending_fatal_container_state_reasons: description: | If the worker pods are in a pending state due to a fatal container 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 67616ea59894d..1dcf813e7e993 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 @@ -53,6 +53,7 @@ from airflow.providers.cncf.kubernetes.kube_config import KubeConfig 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.utils.pod_cleanup import cleanup_kpo_zombie_pods from airflow.providers.cncf.kubernetes.version_compat import AIRFLOW_V_3_0_PLUS from airflow.providers.common.compat.sdk import Stats, conf from airflow.utils.helpers import prune_dict @@ -131,6 +132,7 @@ def __init__(self, *args, **kwargs): self.kube_client: client.CoreV1Api | None = None self.scheduler_job_id: str | None = None self._last_completed_pod_adoption = 0.0 + self._last_kpo_zombie_pod_cleanup = 0.0 self.kubernetes_queue: str | None = None self.task_publish_retries: Counter[TaskInstanceKey] = Counter() self.task_publish_max_retries = self.conf.getint( @@ -407,6 +409,19 @@ def sync(self) -> None: self._last_completed_pod_adoption = now self._adopt_completed_pods(self.kube_client) + if self.kube_config.kpo_zombie_pod_cleanup_enabled: + cleanup_interval = self.kube_config.kpo_zombie_pod_cleanup_interval + if now - self._last_kpo_zombie_pod_cleanup >= cleanup_interval: + self._last_kpo_zombie_pod_cleanup = now + cleanup_kpo_zombie_pods( + list_pods=self._list_pods, + kube_client=self.kube_client, + max_deletes=max(0, self.kube_config.kpo_zombie_pod_cleanup_max_deletes_per_loop), + grace_period_seconds=self.kube_config.kpo_zombie_pod_deletion_grace_period_seconds, + delete_options=self.kube_config.delete_option_kwargs, + kube_client_request_args=self.kube_config.kube_client_request_args, + ) + if self.running: self.log.debug("self.running: %s", self.running) if self.queued_tasks: diff --git a/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/get_provider_info.py b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/get_provider_info.py index df11e51092136..ab955e2e899bd 100644 --- a/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/get_provider_info.py +++ b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/get_provider_info.py @@ -208,6 +208,34 @@ def get_provider_info(): "example": None, "default": "False", }, + "kpo_zombie_pod_cleanup_enabled": { + "description": "If True, periodically delete KubernetesPodOperator pods that no longer have an active matching\ntask instance.\n", + "version_added": None, + "type": "boolean", + "example": None, + "default": "False", + }, + "kpo_zombie_pod_cleanup_interval": { + "description": "How often, in seconds, KubernetesPodOperator zombie pod cleanup runs when enabled.\n", + "version_added": None, + "type": "integer", + "example": None, + "default": "300", + }, + "kpo_zombie_pod_cleanup_max_deletes_per_loop": { + "description": "Maximum number of KubernetesPodOperator zombie pods to delete during one cleanup loop.\n", + "version_added": None, + "type": "integer", + "example": None, + "default": "100", + }, + "kpo_zombie_pod_deletion_grace_period_seconds": { + "description": "Grace period, in seconds, used when deleting KubernetesPodOperator zombie pods.\n", + "version_added": None, + "type": "integer", + "example": None, + "default": "5", + }, "worker_pod_pending_fatal_container_state_reasons": { "description": "If the worker pods are in a pending state due to a fatal container\nstate reasons, then fail the task and delete the worker pod\nif delete_worker_pods is True and delete_worker_pods_on_failure is True.\n", "version_added": "8.1.0", diff --git a/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/kube_config.py b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/kube_config.py index cc9d7fc08fe75..ddd1a570c5d41 100644 --- a/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/kube_config.py +++ b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/kube_config.py @@ -55,6 +55,18 @@ def __init__(self, executor_conf: ExecutorConf | None = None): self.delete_worker_pods_on_failure = self._conf.getboolean( self.kubernetes_section, "delete_worker_pods_on_failure" ) + self.kpo_zombie_pod_cleanup_enabled = self._conf.getboolean( + self.kubernetes_section, "kpo_zombie_pod_cleanup_enabled", fallback=False + ) + self.kpo_zombie_pod_cleanup_interval = self._conf.getint( + self.kubernetes_section, "kpo_zombie_pod_cleanup_interval", fallback=300 + ) + self.kpo_zombie_pod_cleanup_max_deletes_per_loop = self._conf.getint( + self.kubernetes_section, "kpo_zombie_pod_cleanup_max_deletes_per_loop", fallback=100 + ) + self.kpo_zombie_pod_deletion_grace_period_seconds = self._conf.getint( + self.kubernetes_section, "kpo_zombie_pod_deletion_grace_period_seconds", fallback=5 + ) self.worker_pod_pending_fatal_container_state_reasons: list[str] = [] fatal_reasons = self._conf.get( self.kubernetes_section, "worker_pod_pending_fatal_container_state_reasons", fallback="" diff --git a/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/utils/pod_cleanup.py b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/utils/pod_cleanup.py new file mode 100644 index 0000000000000..342f249635246 --- /dev/null +++ b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/utils/pod_cleanup.py @@ -0,0 +1,210 @@ +# Licensed to the Apache Software Foundation (ASF) under one +# or more contributor license agreements. See the NOTICE file +# distributed with this work for additional information +# regarding copyright ownership. The ASF licenses this file +# to you under the Apache License, Version 2.0 (the +# "License"); you may not use this file except in compliance +# with the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, +# software distributed under the License is distributed on an +# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +# KIND, either express or implied. See the License for the +# specific language governing permissions and limitations +# under the License. +"""Clean up KubernetesPodOperator pods that no longer belong to active task instances.""" + +from __future__ import annotations + +import logging +from dataclasses import dataclass +from typing import TYPE_CHECKING + +from kubernetes import client +from kubernetes.client.rest import ApiException +from sqlalchemy import select + +from airflow.models.taskinstance import TaskInstance +from airflow.providers.cncf.kubernetes.pod_generator import make_safe_label_value +from airflow.utils.session import NEW_SESSION, provide_session +from airflow.utils.state import State + +if TYPE_CHECKING: + from collections.abc import Callable, Iterable + + from kubernetes.client import models as k8s + from sqlalchemy.orm import Session + +log = logging.getLogger(__name__) + +KPO_LABEL_SELECTOR = "kubernetes_pod_operator=True" + + +@dataclass(frozen=True) +class _KpoPodKey: + dag_id: str + task_id: str + run_id: str + map_index: int + + +@dataclass(frozen=True) +class _KpoPodRef: + key: _KpoPodKey + try_number: int + pod: k8s.V1Pod + + +@dataclass(frozen=True) +class KpoZombiePodCleanupResult: + """Summary of one KPO zombie pod cleanup pass.""" + + scanned: int + candidates: int + zombies: int + deleted: int + skipped: int + + +def _extract_kpo_pod_ref(pod: k8s.V1Pod) -> _KpoPodRef | None: + labels = pod.metadata.labels or {} + try: + dag_id = labels["dag_id"] + task_id = labels["task_id"] + run_id = labels["run_id"] + try_number = int(labels["try_number"]) + except (KeyError, TypeError, ValueError): + return None + + try: + map_index = int(labels.get("map_index", "-1")) + except (TypeError, ValueError): + return None + + return _KpoPodRef( + key=_KpoPodKey(dag_id=dag_id, task_id=task_id, run_id=run_id, map_index=map_index), + try_number=try_number, + pod=pod, + ) + + +def _build_active_task_instance_try_numbers( + pod_refs: Iterable[_KpoPodRef], *, session: Session +) -> dict[_KpoPodKey, int]: + pod_refs = list(pod_refs) + if not pod_refs: + return {} + + rows = session.execute( + select( + TaskInstance.dag_id, + TaskInstance.task_id, + TaskInstance.run_id, + TaskInstance.map_index, + TaskInstance.try_number, + ) + .where(TaskInstance.state.in_(State.unfinished)) + .where(TaskInstance.map_index.in_({pod_ref.key.map_index for pod_ref in pod_refs})) + ) + + active_try_numbers: dict[_KpoPodKey, int] = {} + for dag_id, task_id, run_id, map_index, try_number in rows: + active_try_numbers[ + _KpoPodKey( + dag_id=make_safe_label_value(dag_id), + task_id=make_safe_label_value(task_id), + run_id=make_safe_label_value(run_id), + map_index=map_index, + ) + ] = try_number + return active_try_numbers + + +def _get_pod_creation_timestamp(pod_ref: _KpoPodRef): + return pod_ref.pod.metadata.creation_timestamp is None, pod_ref.pod.metadata.creation_timestamp + + +def _find_zombie_pods(pod_refs: list[_KpoPodRef], *, session: Session) -> list[k8s.V1Pod]: + active_try_numbers = _build_active_task_instance_try_numbers(pod_refs, session=session) + zombie_refs = [ + pod_ref + for pod_ref in pod_refs + if pod_ref.key not in active_try_numbers or pod_ref.try_number < active_try_numbers[pod_ref.key] + ] + return [pod_ref.pod for pod_ref in sorted(zombie_refs, key=_get_pod_creation_timestamp)] + + +def _delete_pod( + kube_client: client.CoreV1Api, + pod: k8s.V1Pod, + *, + grace_period_seconds: int, + delete_options: dict, + kube_client_request_args: dict, +) -> bool: + pod_name = pod.metadata.name + namespace = pod.metadata.namespace + try: + kube_client.delete_namespaced_pod( + pod_name, + namespace, + body=client.V1DeleteOptions(**{**delete_options, "grace_period_seconds": grace_period_seconds}), + **kube_client_request_args, + ) + except ApiException as e: + if e.status == 404: + return False + log.warning( + "Failed to delete zombie KubernetesPodOperator pod %s in namespace %s: %s", pod_name, namespace, e + ) + return False + return True + + +@provide_session +def cleanup_kpo_zombie_pods( + *, + list_pods: Callable[[dict], list[k8s.V1Pod]], + kube_client: client.CoreV1Api, + max_deletes: int, + grace_period_seconds: int, + delete_options: dict | None = None, + kube_client_request_args: dict | None = None, + session: Session = NEW_SESSION, +) -> KpoZombiePodCleanupResult: + """Delete KubernetesPodOperator pods without an active matching task instance.""" + max_deletes = max(0, max_deletes) + query_kwargs = {"label_selector": KPO_LABEL_SELECTOR} + pods = list_pods(query_kwargs) + pod_refs = [pod_ref for pod in pods if (pod_ref := _extract_kpo_pod_ref(pod))] + zombie_pods = _find_zombie_pods(pod_refs, session=session) + deleted = 0 + for pod in zombie_pods[:max_deletes]: + if _delete_pod( + kube_client, + pod, + grace_period_seconds=grace_period_seconds, + delete_options=delete_options or {}, + kube_client_request_args=kube_client_request_args or {}, + ): + deleted += 1 + + skipped = len(pods) - len(pod_refs) + max(0, len(zombie_pods) - max_deletes) + result = KpoZombiePodCleanupResult( + scanned=len(pods), + candidates=len(pod_refs), + zombies=len(zombie_pods), + deleted=deleted, + skipped=skipped, + ) + log.info( + "KPO zombie pod cleanup scanned=%s candidates=%s zombies=%s deleted=%s skipped=%s", + result.scanned, + result.candidates, + result.zombies, + result.deleted, + result.skipped, + ) + return result 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 a71101dc55515..69481b6487665 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 @@ -2388,6 +2388,56 @@ def test_alive_other_scheduler_job_ids_does_not_detach_caller_session(self, sess "_alive_other_scheduler_job_ids closed/detached the caller's scoped session" ) + @mock.patch("airflow.providers.cncf.kubernetes.executors.kubernetes_executor.cleanup_kpo_zombie_pods") + @mock.patch( + "airflow.providers.cncf.kubernetes.executors.kubernetes_executor.KubernetesExecutor._adopt_completed_pods" + ) + @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_skips_kpo_zombie_pod_cleanup_when_disabled( + self, mock_get_kube_client, mock_kubernetes_job_watcher, mock_adopt_completed_pods, mock_cleanup + ): + executor = self.kubernetes_executor + executor.start() + try: + executor.kube_config.kpo_zombie_pod_cleanup_enabled = False + + executor.sync() + + mock_cleanup.assert_not_called() + finally: + executor.end() + + @mock.patch("airflow.providers.cncf.kubernetes.executors.kubernetes_executor.cleanup_kpo_zombie_pods") + @mock.patch( + "airflow.providers.cncf.kubernetes.executors.kubernetes_executor.KubernetesExecutor._adopt_completed_pods" + ) + @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_runs_kpo_zombie_pod_cleanup_when_enabled( + self, mock_get_kube_client, mock_kubernetes_job_watcher, mock_adopt_completed_pods, mock_cleanup + ): + executor = self.kubernetes_executor + executor.start() + try: + executor.kube_config.kpo_zombie_pod_cleanup_enabled = True + executor.kube_config.kpo_zombie_pod_cleanup_interval = 0 + executor.kube_config.kpo_zombie_pod_cleanup_max_deletes_per_loop = 100 + executor.kube_config.kpo_zombie_pod_deletion_grace_period_seconds = 5 + + executor.sync() + + mock_cleanup.assert_called_once_with( + list_pods=executor._list_pods, + kube_client=executor.kube_client, + max_deletes=100, + grace_period_seconds=5, + delete_options=executor.kube_config.delete_option_kwargs, + kube_client_request_args=executor.kube_config.kube_client_request_args, + ) + finally: + executor.end() + @pytest.mark.db_test @mock.patch("airflow.providers.cncf.kubernetes.executors.kubernetes_executor_utils.KubernetesJobWatcher") @mock.patch("airflow.providers.cncf.kubernetes.kube_client.get_kube_client") diff --git a/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/utils/test_pod_cleanup.py b/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/utils/test_pod_cleanup.py new file mode 100644 index 0000000000000..6d69aeb9c7a70 --- /dev/null +++ b/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/utils/test_pod_cleanup.py @@ -0,0 +1,156 @@ +# Licensed to the Apache Software Foundation (ASF) under one +# or more contributor license agreements. See the NOTICE file +# distributed with this work for additional information +# regarding copyright ownership. The ASF licenses this file +# to you under the Apache License, Version 2.0 (the +# "License"); you may not use this file except in compliance +# with the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, +# software distributed under the License is distributed on an +# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +# KIND, either express or implied. See the License for the +# specific language governing permissions and limitations +# under the License. +from __future__ import annotations + +from datetime import datetime, timezone + +from kubernetes.client import models as k8s + +from airflow.providers.cncf.kubernetes.pod_generator import make_safe_label_value +from airflow.providers.cncf.kubernetes.utils.pod_cleanup import cleanup_kpo_zombie_pods + + +class FakeSession: + def __init__(self, rows): + self.rows = rows + + def execute(self, _statement): + return self.rows + + +class FakeKubeClient: + def __init__(self): + self.deleted = [] + + def delete_namespaced_pod(self, name, namespace, body, **kwargs): + self.deleted.append((name, namespace, body, kwargs)) + + +def make_kpo_pod( + name, + *, + dag_id="dag", + task_id="task", + run_id="run", + try_number="1", + map_index=None, + created_at=None, +): + labels = { + "dag_id": make_safe_label_value(dag_id), + "task_id": make_safe_label_value(task_id), + "run_id": make_safe_label_value(run_id), + "try_number": try_number, + "kubernetes_pod_operator": "True", + } + if map_index is not None: + labels["map_index"] = str(map_index) + return k8s.V1Pod( + metadata=k8s.V1ObjectMeta( + name=name, namespace="default", labels=labels, creation_timestamp=created_at + ) + ) + + +def make_task_instance_row(*, dag_id="dag", task_id="task", run_id="run", map_index=-1, try_number=1): + return dag_id, task_id, run_id, map_index, try_number + + +def run_cleanup(pods, rows, *, max_deletes=100): + kube_client = FakeKubeClient() + result = cleanup_kpo_zombie_pods( + list_pods=lambda _query_kwargs: pods, + kube_client=kube_client, + max_deletes=max_deletes, + grace_period_seconds=5, + delete_options={"propagation_policy": "Foreground"}, + kube_client_request_args={"request_timeout": 10}, + session=FakeSession(rows), + ) + return result, kube_client + + +def test_cleanup_deletes_pod_without_active_task_instance(): + pod = make_kpo_pod("zombie") + + result, kube_client = run_cleanup([pod], []) + + assert result.scanned == 1 + assert result.zombies == 1 + assert result.deleted == 1 + name, namespace, body, kwargs = kube_client.deleted[0] + assert name == "zombie" + assert namespace == "default" + assert body.grace_period_seconds == 5 + assert body.propagation_policy == "Foreground" + assert kwargs == {"request_timeout": 10} + + +def test_cleanup_keeps_pod_with_active_matching_task_instance(): + pod = make_kpo_pod("active") + + result, kube_client = run_cleanup([pod], [make_task_instance_row()]) + + assert result.zombies == 0 + assert result.deleted == 0 + assert kube_client.deleted == [] + + +def test_cleanup_deletes_old_retry_pod(): + pod = make_kpo_pod("old-try", try_number="1") + + result, kube_client = run_cleanup([pod], [make_task_instance_row(try_number=2)]) + + assert result.zombies == 1 + assert result.deleted == 1 + assert kube_client.deleted[0][0] == "old-try" + + +def test_cleanup_normalizes_active_task_instance_labels_before_matching(): + long_run_id = "scheduled__" + "a" * 100 + ":with:colons" + pod = make_kpo_pod("active-long-run", run_id=long_run_id) + + result, kube_client = run_cleanup([pod], [make_task_instance_row(run_id=long_run_id)]) + + assert result.zombies == 0 + assert result.deleted == 0 + assert kube_client.deleted == [] + + +def test_cleanup_skips_pods_with_incomplete_labels(): + pod = make_kpo_pod("missing-try") + del pod.metadata.labels["try_number"] + + result, kube_client = run_cleanup([pod], []) + + assert result.scanned == 1 + assert result.candidates == 0 + assert result.zombies == 0 + assert result.skipped == 1 + assert kube_client.deleted == [] + + +def test_cleanup_respects_max_deletes_and_deletes_oldest_pods_first(): + newer = make_kpo_pod("newer", created_at=datetime(2024, 1, 2, tzinfo=timezone.utc)) + older = make_kpo_pod("older", created_at=datetime(2024, 1, 1, tzinfo=timezone.utc)) + + result, kube_client = run_cleanup([newer, older], [], max_deletes=1) + + assert result.zombies == 2 + assert result.deleted == 1 + assert result.skipped == 1 + assert kube_client.deleted[0][0] == "older"