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 @@ -1619,6 +1619,7 @@ StatsD
statsd
stderr
stdin
stdlib
stdout
StorageClass
storages
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,7 @@
from datetime import timedelta
from typing import TYPE_CHECKING, Any, cast

import tenacity
from botocore.exceptions import ClientError, WaiterError

from airflow.exceptions import AirflowProviderDeprecationWarning
Expand All @@ -42,7 +43,10 @@
EksDeleteNodegroupTrigger,
EksPodTrigger,
)
from airflow.providers.amazon.aws.utils import validate_execute_complete_event
from airflow.providers.amazon.aws.utils import (
build_resource_in_use_retry_args,
validate_execute_complete_event,
)
from airflow.providers.amazon.aws.utils.mixins import aws_template_fields
from airflow.providers.amazon.aws.utils.waiter_with_logging import wait
from airflow.providers.cncf.kubernetes.utils.pod_manager import OnFinishAction
Expand Down Expand Up @@ -139,10 +143,12 @@ def _create_compute(
# delete_nodegroup_on_failure defaults to True to prevent orphaned nodegroups.
if delete_nodegroup_on_failure:
try:
eks_hook.delete_nodegroup(
clusterName=cluster_name,
nodegroupName=nodegroup_name,
)
for attempt in tenacity.Retrying(**build_resource_in_use_retry_args(log)):
with attempt:
eks_hook.delete_nodegroup(
clusterName=cluster_name,
nodegroupName=nodegroup_name,
)
log.info(
"Issued delete request for nodegroup '%s' in cluster '%s' after failure.",
nodegroup_name,
Expand Down Expand Up @@ -809,7 +815,9 @@ def execute(self, context: Context):
self.delete_any_nodegroups()
self.delete_any_fargate_profiles()

self.hook.delete_cluster(name=self.cluster_name)
for attempt in tenacity.Retrying(**build_resource_in_use_retry_args(self.log)):
with attempt:
self.hook.delete_cluster(name=self.cluster_name)

if self.wait_for_completion:
self.log.info("Waiting for cluster to delete. This will take some time.")
Expand All @@ -825,8 +833,11 @@ def delete_any_nodegroups(self) -> None:
nodegroups = self.hook.list_nodegroups(clusterName=self.cluster_name)
if nodegroups:
self.log.info(CAN_NOT_DELETE_MSG.format(compute=NODEGROUP_FULL_NAME, count=len(nodegroups)))
retry_args = build_resource_in_use_retry_args(self.log)
for group in nodegroups:
self.hook.delete_nodegroup(clusterName=self.cluster_name, nodegroupName=group)
for attempt in tenacity.Retrying(**retry_args):
with attempt:
self.hook.delete_nodegroup(clusterName=self.cluster_name, nodegroupName=group)
# Note this is a custom waiter so we're using hook.get_waiter(), not hook.conn.get_waiter().
self.log.info("Waiting for all nodegroups to delete. This will take some time.")
self.hook.get_waiter("all_nodegroups_deleted").wait(clusterName=self.cluster_name)
Expand All @@ -843,11 +854,16 @@ def delete_any_fargate_profiles(self) -> None:
if fargate_profiles:
self.log.info(CAN_NOT_DELETE_MSG.format(compute=FARGATE_FULL_NAME, count=len(fargate_profiles)))
self.log.info("Waiting for Fargate profiles to delete. This will take some time.")
retry_args = build_resource_in_use_retry_args(self.log)
for profile in fargate_profiles:
# The API will return a (cluster) ResourceInUseException if you try
# to delete Fargate profiles in parallel the way we can with nodegroups,
# so each must be deleted sequentially
self.hook.delete_fargate_profile(clusterName=self.cluster_name, fargateProfileName=profile)
for attempt in tenacity.Retrying(**retry_args):
with attempt:
self.hook.delete_fargate_profile(
clusterName=self.cluster_name, fargateProfileName=profile
)
self.hook.conn.get_waiter("fargate_profile_deleted").wait(
clusterName=self.cluster_name, fargateProfileName=profile
)
Expand Down Expand Up @@ -921,7 +937,9 @@ def __init__(
super().__init__(**kwargs)

def execute(self, context: Context):
self.hook.delete_nodegroup(clusterName=self.cluster_name, nodegroupName=self.nodegroup_name)
for attempt in tenacity.Retrying(**build_resource_in_use_retry_args(self.log)):
with attempt:
self.hook.delete_nodegroup(clusterName=self.cluster_name, nodegroupName=self.nodegroup_name)
if self.deferrable:
self.defer(
trigger=EksDeleteNodegroupTrigger(
Expand Down Expand Up @@ -1009,9 +1027,11 @@ def __init__(
super().__init__(**kwargs)

def execute(self, context: Context):
self.hook.delete_fargate_profile(
clusterName=self.cluster_name, fargateProfileName=self.fargate_profile_name
)
for attempt in tenacity.Retrying(**build_resource_in_use_retry_args(self.log)):
with attempt:
self.hook.delete_fargate_profile(
clusterName=self.cluster_name, fargateProfileName=self.fargate_profile_name
)
if self.deferrable:
self.defer(
trigger=EksDeleteFargateProfileTrigger(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,10 +19,12 @@
import datetime
from typing import TYPE_CHECKING, Any

import tenacity
from botocore.exceptions import ClientError

from airflow.providers.amazon.aws.hooks.eks import EksHook
from airflow.providers.amazon.aws.triggers.base import AwsBaseWaiterTrigger
from airflow.providers.amazon.aws.utils import build_resource_in_use_retry_args
from airflow.providers.amazon.aws.utils.waiter_with_logging import async_wait
from airflow.providers.cncf.kubernetes.triggers.pod import KubernetesPodTrigger
from airflow.providers.common.compat.sdk import AirflowException
Expand Down Expand Up @@ -275,13 +277,15 @@ async def run(self):
if self.force_delete_compute:
await self.delete_any_nodegroups(client=client)
await self.delete_any_fargate_profiles(client=client)
try:
await client.delete_cluster(name=self.cluster_name)
except ClientError as ex:
if ex.response.get("Error").get("Code") == "ResourceNotFoundException":
pass
else:
raise
async for attempt in tenacity.AsyncRetrying(**build_resource_in_use_retry_args(self.log)):
with attempt:
try:
await client.delete_cluster(name=self.cluster_name)
except ClientError as ex:
# The cluster is already gone — nothing to wait on, so stop retrying.
if ex.response.get("Error", {}).get("Code") == "ResourceNotFoundException":
break
raise
await async_wait(
waiter=waiter,
waiter_delay=int(self.waiter_delay),
Expand All @@ -305,8 +309,11 @@ async def delete_any_nodegroups(self, client) -> None:
if nodegroups.get("nodegroups", None):
self.log.info("Deleting nodegroups")
waiter = self.hook().get_waiter("all_nodegroups_deleted", deferrable=True, client=client)
retry_args = build_resource_in_use_retry_args(self.log)
for group in nodegroups["nodegroups"]:
await client.delete_nodegroup(clusterName=self.cluster_name, nodegroupName=group)
async for attempt in tenacity.AsyncRetrying(**retry_args):
with attempt:
await client.delete_nodegroup(clusterName=self.cluster_name, nodegroupName=group)
await async_wait(
waiter=waiter,
waiter_delay=int(self.waiter_delay),
Expand All @@ -330,8 +337,13 @@ async def delete_any_fargate_profiles(self, client) -> None:
fargate_profiles = await client.list_fargate_profiles(clusterName=self.cluster_name)
if fargate_profiles.get("fargateProfileNames"):
self.log.info("Waiting for Fargate profiles to delete. This will take some time.")
retry_args = build_resource_in_use_retry_args(self.log)
for profile in fargate_profiles["fargateProfileNames"]:
await client.delete_fargate_profile(clusterName=self.cluster_name, fargateProfileName=profile)
async for attempt in tenacity.AsyncRetrying(**retry_args):
with attempt:
await client.delete_fargate_profile(
clusterName=self.cluster_name, fargateProfileName=profile
)
await async_wait(
waiter=client.get_waiter("fargate_profile_deleted"),
waiter_delay=int(self.waiter_delay),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -22,14 +22,54 @@
from datetime import datetime, timezone
from enum import Enum
from importlib import metadata
from typing import Any
from typing import TYPE_CHECKING, Any

import tenacity
from botocore.exceptions import ClientError

from airflow.providers.common.compat.sdk import AirflowException
from airflow.utils.helpers import prune_dict
from airflow.version import version

if TYPE_CHECKING:
from airflow.sdk.types import Logger

log = logging.getLogger(__name__)

# AWS briefly rejects a delete call with ResourceInUseException while the target resource is still
# settling from a prior operation (e.g. EKS finalizing a nodegroup removal before the cluster can be
# deleted). Retry with exponential backoff (1s, 2s, 4s, ... capped at RESOURCE_IN_USE_RETRY_MAX_WAIT
# per wait) until RESOURCE_IN_USE_RETRY_TIMEOUT elapses, then give up and re-raise. This rides out the
# settling window without hanging a genuinely wedged resource for long.
RESOURCE_IN_USE_RETRY_TIMEOUT = 300
RESOURCE_IN_USE_RETRY_MAX_WAIT = 60


def is_resource_in_use_error(exception: BaseException) -> bool:
"""Return True if the exception is a transient AWS ``ResourceInUseException``."""
return (
isinstance(exception, ClientError)
and exception.response.get("Error", {}).get("Code") == "ResourceInUseException"
)


def build_resource_in_use_retry_args(logger: Logger | logging.Logger) -> dict[str, Any]:
"""
Build tenacity arguments for retrying a call on a transient ``ResourceInUseException``.

Shared by synchronous operators (``tenacity.Retrying``) and deferrable triggers
(``tenacity.AsyncRetrying``) so both back off identically. ``reraise=True`` keeps the
original error as the task failure once the retry timeout is exhausted. Accepts either an
Airflow structlog logger (``self.log``) or a stdlib ``logging.Logger`` (module-level helpers).
"""
return {
"retry": tenacity.retry_if_exception(is_resource_in_use_error),
"wait": tenacity.wait_exponential(max=RESOURCE_IN_USE_RETRY_MAX_WAIT),
"stop": tenacity.stop_after_delay(RESOURCE_IN_USE_RETRY_TIMEOUT),
"before_sleep": tenacity.before_sleep_log(logger, logging.WARNING),
"reraise": True,
}


def trim_none_values(obj: dict):
return prune_dict(obj)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,8 +18,6 @@

from datetime import datetime

from pendulum import duration

from airflow.providers.amazon.aws.hooks.eks import ClusterStates, FargateProfileStates
from airflow.providers.amazon.aws.operators.eks import (
EksCreateClusterOperator,
Expand Down Expand Up @@ -132,9 +130,6 @@
trigger_rule=TriggerRule.ALL_DONE,
cluster_name=cluster_name,
force_delete_compute=True,
retries=4,
retry_delay=duration(seconds=30),
retry_exponential_backoff=True,
)

await_delete_cluster = EksClusterStateSensor(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,8 +18,6 @@

from datetime import datetime

from pendulum import duration

from airflow.providers.amazon.aws.hooks.eks import ClusterStates, FargateProfileStates
from airflow.providers.amazon.aws.operators.eks import (
EksCreateClusterOperator,
Expand Down Expand Up @@ -146,9 +144,6 @@
task_id="delete_eks_fargate_profile",
cluster_name=cluster_name,
fargate_profile_name=fargate_profile_name,
retries=4,
retry_delay=duration(seconds=30),
retry_exponential_backoff=True,
)
# [END howto_operator_eks_delete_fargate_profile]
delete_fargate_profile.trigger_rule = TriggerRule.ALL_DONE
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,6 @@
from datetime import datetime

import boto3
from pendulum import duration

from airflow.providers.amazon.aws.hooks.eks import ClusterStates, NodegroupStates
from airflow.providers.amazon.aws.operators.eks import (
Expand Down Expand Up @@ -149,9 +148,6 @@ def delete_launch_template(template_name: str):
task_id="delete_nodegroup_and_cluster",
cluster_name=cluster_name,
force_delete_compute=True,
retries=4,
retry_delay=duration(seconds=30),
retry_exponential_backoff=True,
)
# [END howto_operator_eks_force_delete_cluster]
delete_nodegroup_and_cluster.trigger_rule = TriggerRule.ALL_DONE
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,6 @@
from datetime import datetime

import boto3
from pendulum import duration

from airflow.providers.amazon.aws.hooks.eks import ClusterStates, NodegroupStates
from airflow.providers.amazon.aws.operators.eks import (
Expand Down Expand Up @@ -172,9 +171,6 @@ def delete_launch_template(template_name: str):
task_id="delete_nodegroup",
cluster_name=cluster_name,
nodegroup_name=nodegroup_name,
retries=4,
retry_delay=duration(seconds=30),
retry_exponential_backoff=True,
)
# [END howto_operator_eks_delete_nodegroup]
delete_nodegroup.trigger_rule = TriggerRule.ALL_DONE
Expand Down
Loading