Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions docs/spelling_wordlist.txt
Original file line number Diff line number Diff line change
Expand Up @@ -959,6 +959,7 @@ kinesis
kinit
kms
knownHosts
kpo
krb
Kube
kube
Expand Down
1 change: 1 addition & 0 deletions providers/cncf/kubernetes/docs/changelog.rst
Original file line number Diff line number Diff line change
Expand Up @@ -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)``

Expand Down
29 changes: 29 additions & 0 deletions providers/cncf/kubernetes/provider.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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(
Expand Down Expand Up @@ -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:
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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=""
Expand Down
Original file line number Diff line number Diff line change
@@ -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
Loading