Skip to content
2 changes: 2 additions & 0 deletions airflow-core/tests/unit/always/test_project_structure.py
Original file line number Diff line number Diff line change
Expand Up @@ -498,6 +498,8 @@ class TestAmazonProviderProjectStructure(ExampleCoverageTest):

BASE_CLASSES = {
"airflow.providers.amazon.aws.operators.base_aws.AwsBaseOperator",
# Compatibility fallback for EksPodExecOperator, not a standalone Amazon operator.
"airflow.providers.amazon.aws.operators.eks.KubernetesPodExecOperator",
"airflow.providers.amazon.aws.operators.glue_crawler._GlueCrawlerBaseOperator",
"airflow.providers.amazon.aws.operators.rds.RdsBaseOperator",
"airflow.providers.amazon.aws.operators.sagemaker.SageMakerBaseOperator",
Expand Down
30 changes: 30 additions & 0 deletions providers/amazon/docs/operators/eks.rst
Original file line number Diff line number Diff line change
Expand Up @@ -205,6 +205,36 @@ Note: An Amazon EKS Cluster with underlying compute infrastructure is required.
:start-after: [START howto_operator_eks_pod_operator]
:end-before: [END howto_operator_eks_pod_operator]

.. _howto/operator:EksPodExecOperator:

Execute a command in an existing Pod on Amazon EKS
==================================================

To execute a command in a running container without managing the Pod lifecycle, use
:class:`~airflow.providers.amazon.aws.operators.eks.EksPodExecOperator`.

This operator requires ``apache-airflow-providers-cncf-kubernetes>=10.22.0``.
Existing EKS operators remain available with older supported versions of the Kubernetes provider.

As with ``EksPodOperator``, ``kubernetes_conn_id`` defaults to ``kubernetes_default`` and can be
set to another Kubernetes connection. If the default connection contains ``kube_config`` or
``cluster_context``, use a separate connection without those settings, since EKS generates its own kubeconfig.

The Pod must already exist and be running. The operator streams command output, waits for the exit code,
and does not create, restart, or delete the Pod.

The AWS identity must have permission to call ``eks:DescribeCluster`` and be authorized to access the
EKS cluster. Kubernetes RBAC must allow ``get`` on ``pods`` and ``pods/exec``.

See :class:`~airflow.providers.cncf.kubernetes.operators.pod_exec.KubernetesPodExecOperator`
Comment thread
AlejandroMorgante marked this conversation as resolved.
for command, output, XCom, and retry behavior.

.. exampleinclude:: /../../amazon/tests/system/amazon/aws/example_eks_with_nodegroups.py
:language: python
:dedent: 4
:start-after: [START howto_operator_eks_pod_exec]
:end-before: [END howto_operator_eks_pod_exec]

Sensors
-------

Expand Down
132 changes: 130 additions & 2 deletions providers/amazon/src/airflow/providers/amazon/aws/operators/eks.py
Original file line number Diff line number Diff line change
Expand Up @@ -47,10 +47,28 @@
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.mixins import AwsHookParams, aws_template_fields
from airflow.providers.amazon.aws.utils.waiter_with_logging import wait
from airflow.providers.cncf.kubernetes.utils.pod_manager import OnFinishAction
from airflow.providers.common.compat.sdk import AirflowException, conf
from airflow.providers.common.compat.sdk import (
AirflowException,
AirflowOptionalProviderFeatureException,
BaseOperator,
conf,
)

try:
from airflow.providers.cncf.kubernetes.operators.pod_exec import KubernetesPodExecOperator
except ImportError:

class KubernetesPodExecOperator(BaseOperator): # type: ignore[no-redef]
"""Keep existing EKS operators importable with older Kubernetes providers."""

def __init__(self, **kwargs):
raise AirflowOptionalProviderFeatureException(
"EksPodExecOperator requires apache-airflow-providers-cncf-kubernetes>=10.22.0."
)

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I based this on the existing backward-compatible import fallback in the EKS module

try:
from airflow.providers.cncf.kubernetes.operators.pod import KubernetesPodOperator
except ImportError:
# preserve backward compatibility for older versions of cncf.kubernetes provider, remove this when minimum cncf.kubernetes provider is 10.0
from airflow.providers.cncf.kubernetes.operators.kubernetes_pod import ( # type: ignore[no-redef]
KubernetesPodOperator,
)

The difference is that KubernetesPodExecOperator has no equivalent in older versions of the cncf-kubernetes provider. This fallback keeps existing EKS operators importable and only raises an error when someone tries to instantiate EksPodExecOperator without the required provider version.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Sorry. I just saw that is just a few lines down in the existing file. My mistake.



try:
from airflow.providers.cncf.kubernetes.operators.pod import KubernetesPodOperator
Expand Down Expand Up @@ -1368,3 +1386,113 @@ def _refresh_cached_properties(self) -> None:
self.log.exception("Failed to refresh AWS credentials.")
raise
super()._refresh_cached_properties()


class EksPodExecOperator(KubernetesPodExecOperator):
"""
Execute a command in a running container of an existing Pod on Amazon EKS.

The operator authenticates with Amazon EKS and delegates command execution to
:class:`~airflow.providers.cncf.kubernetes.operators.pod_exec.KubernetesPodExecOperator`.
It does not create, restart, or delete the target Pod.

.. seealso::
For more information on how to use this operator, take a look at the guide:
:ref:`howto/operator:EksPodExecOperator`

:param cluster_name: The name of the Amazon EKS Cluster containing the Pod. (templated)
:param pod_name: Name of the existing Kubernetes Pod. (templated)
:param command: Command and arguments to execute in the container. (templated)
:param namespace: Namespace containing the Pod. Defaults to ``default``. (templated)
:param container_name: Name of the container in which to execute the command. When omitted, the
``kubectl.kubernetes.io/default-container`` annotation or the first container is used.
Defaults to ``None``. (templated)
:param aws_conn_id: The Airflow connection used for AWS credentials. (templated)
Defaults to ``aws_default``. If this is ``None`` or empty, the default boto3 credential
strategy is used without an Airflow connection lookup.
:param region_name: AWS region containing the Amazon EKS Cluster. (templated)
Defaults to ``None``, which uses the region from the AWS connection when available and
otherwise falls back to the default boto3 region strategy.
:param verify: Whether to verify SSL certificates, or the path to a CA bundle. Defaults to
``None``, which uses the value from the AWS connection when available. (templated)
:param botocore_config: Configuration dictionary for the botocore client. Defaults to ``None``,
which uses ``config_kwargs`` from the AWS connection when available.
:param kubernetes_conn_id: Kubernetes connection used for additional client configuration.
Defaults to ``kubernetes_default``, as with ``EksPodOperator``. (templated)
:param do_xcom_push: Return standard output through XCom when ``True``. Defaults to ``False``.
:param max_xcom_output_size: Maximum UTF-8 byte size retained for XCom. Defaults to 49,344 bytes.
"""

template_fields: Sequence[str] = aws_template_fields(
"cluster_name",
*(
field
for field in KubernetesPodExecOperator.template_fields
if field not in {"cluster_context", "config_file"}
),
)

def __init__(
self,
*,
cluster_name: str,
pod_name: str,
command: Sequence[str],
namespace: str = DEFAULT_NAMESPACE_NAME,
container_name: str | None = None,
aws_conn_id: str | None = DEFAULT_CONN_ID,
region_name: str | None = None,
verify: bool | str | None = None,
botocore_config: dict | None = None,
**kwargs,
) -> None:
hook_params = AwsHookParams.from_constructor(
aws_conn_id, region_name, verify, botocore_config, additional_params=kwargs
)
super().__init__(
pod_name=pod_name,
command=command,
namespace=namespace,
container_name=container_name,
in_cluster=False,
cluster_context=None,
config_file=None,
**kwargs,
)
self.cluster_name = cluster_name
self.aws_conn_id = hook_params.aws_conn_id
self.region_name = hook_params.region_name
self.verify = hook_params.verify
self.botocore_config = hook_params.botocore_config

def execute(self, context: Context) -> str | None:
eks_hook = EksHook(
aws_conn_id=self.aws_conn_id,
region_name=self.region_name,
verify=self.verify,
config=self.botocore_config,
)
credentials = eks_hook.get_session().get_credentials()
if credentials is None:
raise RuntimeError(
"Unable to retrieve AWS credentials. Credentials may have expired or not been configured. "
"Please check your AWS connection configuration."
)
frozen_credentials = credentials.get_frozen_credentials()
with eks_hook._secure_credential_context(
frozen_credentials.access_key,
frozen_credentials.secret_key,
frozen_credentials.token,
) as credentials_file:
with eks_hook.generate_config_file(
eks_cluster_name=self.cluster_name,
pod_namespace=self.namespace,
credentials_file=credentials_file,
) as config_file:
self.config_file = config_file
try:
return super().execute(context)
finally:
self.config_file = None
Comment thread
AlejandroMorgante marked this conversation as resolved.
self.__dict__.pop("client", None)
self.__dict__.pop("hook", None)
Original file line number Diff line number Diff line change
Expand Up @@ -16,19 +16,26 @@
# under the License.
from __future__ import annotations

import asyncio
from contextlib import contextmanager
from datetime import datetime
from typing import TYPE_CHECKING

import boto3
from kubernetes.client import V1Container, V1ObjectMeta, V1Pod, V1PodSpec

from airflow.providers.amazon.aws.hooks.eks import ClusterStates, NodegroupStates
from airflow.providers.amazon.aws.hooks.eks import ClusterStates, EksHook, NodegroupStates
from airflow.providers.amazon.aws.operators.eks import (
EksCreateClusterOperator,
EksCreateNodegroupOperator,
EksDeleteClusterOperator,
EksDeleteNodegroupOperator,
EksPodExecOperator,
EksPodOperator,
)
from airflow.providers.amazon.aws.sensors.eks import EksClusterStateSensor, EksNodegroupStateSensor
from airflow.providers.cncf.kubernetes.hooks.kubernetes import KubernetesHook
from airflow.providers.cncf.kubernetes.utils.pod_manager import PodManager

from tests_common.test_utils.version_compat import AIRFLOW_V_3_0_PLUS

Expand All @@ -49,7 +56,13 @@
from system.amazon.aws.utils import ENV_ID_KEY, SystemTestContextBuilder
from system.amazon.aws.utils.k8s import get_describe_pod_operator

if TYPE_CHECKING:
from collections.abc import Generator

from kubernetes.client import CoreV1Api

DAG_ID = "example_eks_with_nodegroups"
EXPECTED_EXEC_OUTPUT = "command executed in existing EKS pod"

# Externally fetched variables:
ROLE_ARN_KEY = "ROLE_ARN"
Expand All @@ -76,6 +89,50 @@ def delete_launch_template(template_name: str):
boto3.client("ec2").delete_launch_template(LaunchTemplateName=template_name)


@contextmanager
def get_eks_kubernetes_client(cluster_name: str) -> Generator[CoreV1Api, None, None]:
eks_hook = EksHook()
credentials = eks_hook.get_session().get_credentials()
if credentials is None:
raise RuntimeError("Unable to retrieve AWS credentials for the EKS system test.")
frozen_credentials = credentials.get_frozen_credentials()
with eks_hook._secure_credential_context(
frozen_credentials.access_key,
frozen_credentials.secret_key,
frozen_credentials.token,
) as credentials_file:
with eks_hook.generate_config_file(cluster_name, "default", credentials_file) as config_file:
yield KubernetesHook(kubernetes_conn_id=None, config_file=config_file).core_v1_client


@task
def create_exec_pod(cluster_name: str, pod_name: str) -> None:
pod = V1Pod(
metadata=V1ObjectMeta(name=pod_name, namespace="default"),
spec=V1PodSpec(
containers=[V1Container(name="main", image="busybox:1.38.0", command=["sleep", "3600"])],
restart_policy="Never",
),
)
with get_eks_kubernetes_client(cluster_name) as kube_client:
pod_manager = PodManager(kube_client=kube_client)
created_pod = pod_manager.create_pod(pod)
asyncio.run(pod_manager.await_pod_start(created_pod))


@task(trigger_rule=TriggerRule.ALL_DONE)
def delete_exec_pod(cluster_name: str, pod_name: str) -> None:
pod = V1Pod(metadata=V1ObjectMeta(name=pod_name, namespace="default"))
with get_eks_kubernetes_client(cluster_name) as kube_client:
PodManager(kube_client=kube_client).delete_pod(pod)


@task
def verify_exec_output(output: str) -> None:
if output != EXPECTED_EXEC_OUTPUT:
raise ValueError(f"Unexpected command output: {output!r}")


with DAG(
dag_id=DAG_ID,
schedule="@once",
Expand All @@ -88,6 +145,7 @@ def delete_launch_template(template_name: str):
cluster_name = f"{env_id}-cluster"
nodegroup_name = f"{env_id}-nodegroup"
launch_template_name = f"{env_id}-launch-template"
exec_pod_name = f"{env_id}-exec-pod"

# [START howto_operator_eks_create_cluster]
# Create an Amazon EKS Cluster control plane without attaching compute service.
Expand Down Expand Up @@ -151,6 +209,22 @@ def delete_launch_template(template_name: str):
# it is cleaned anyway with the cluster later on.
start_pod.is_delete_operator_pod = False

create_exec_pod_task = create_exec_pod(cluster_name, exec_pod_name)

# [START howto_operator_eks_pod_exec]
run_command = EksPodExecOperator(
task_id="run_command_in_existing_pod",
cluster_name=cluster_name,
pod_name=exec_pod_name,
command=["sh", "-c", f"printf '{EXPECTED_EXEC_OUTPUT}'"],
do_xcom_push=True,
)
# [END howto_operator_eks_pod_exec]

exec_output_is_valid = verify_exec_output(run_command.output)

delete_exec_pod_task = delete_exec_pod(cluster_name, exec_pod_name)

describe_pod = get_describe_pod_operator(
cluster_name, pod_name="{{ ti.xcom_pull(key='pod_name', task_ids='run_pod') }}"
)
Expand Down Expand Up @@ -211,15 +285,31 @@ def delete_launch_template(template_name: str):
# TEST SETUP
test_context,
create_launch_template(launch_template_name),
# TEST BODY
create_cluster,
await_create_cluster,
create_nodegroup,
await_create_nodegroup,
)
chain(
# TEST BODY: EksPodOperator
await_create_nodegroup,
start_pod,
# TEST TEARDOWN
describe_pod,
await_nodegroup_stable,
)
chain(
# TEST BODY: EksPodExecOperator
await_create_nodegroup,
create_exec_pod_task,
run_command,
exec_output_is_valid,
# TEST TEARDOWN
delete_exec_pod_task,
await_nodegroup_stable,
)
chain(
# TEST TEARDOWN
await_nodegroup_stable,
delete_nodegroup, # part of the test AND teardown
await_delete_nodegroup,
await_cluster_stable,
Expand Down
Loading