eladkal commented on code in PR #74399: URL: https://github.com/apache/airflow/pull/74399#discussion_r4218475819
########## providers/amazon/docs/executors/eks-executor.rst: ########## @@ -0,0 +1,232 @@ + .. 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. + +.. |executorName| replace:: EKS + +.. _eks_executor: + +================ +AWS EKS Executor +================ + +The EKS executor runs each Airflow task in its own pod on an Amazon EKS cluster. + +This executor extends the Kubernetes executor that ships in the ``cncf.kubernetes`` +provider. Pod scheduling, pod templates, per-task pod overrides, and log handling all +behave the same way they do under +:doc:`apache-airflow-providers-cncf-kubernetes:kubernetes_executor`. The part this executor +adds is authentication. It builds the Kubernetes client for your EKS cluster from an +Airflow AWS connection, so the scheduler does not need a kubeconfig file on disk and +you do not have to refresh cluster credentials yourself. + +Because the executor inherits its pod behavior, everything written about the +Kubernetes executor still applies, including the requirement for a database backend +other than SQLite. + +For a quick start guide please see :ref:`here <eks_setup_guide>`. + +Requirements +------------ + +The executor needs ``apache-airflow-providers-cncf-kubernetes`` version 10.24.0 or +newer, the release that added the ``client_factory`` setting the executor is built on. +Install the ``cncf.kubernetes`` extra of the Amazon provider together with that +version: + +.. code-block:: bash + + pip install 'apache-airflow-providers-amazon[cncf.kubernetes]' \ + 'apache-airflow-providers-cncf-kubernetes>=10.24.0' + +The extra by itself declares a much older floor, because the Amazon provider's other +Kubernetes integrations still work with it. A fresh install resolves to the newest +``cncf.kubernetes`` regardless, so naming the version matters when an older one is +already pinned in your environment. If the installed version is too old, the executor +raises an error when it starts up rather than falling back to a client it cannot +authenticate with. + +The AWS credentials the executor uses must be allowed to call ``eks:DescribeCluster`` +on the cluster, and the IAM principal behind those credentials must be granted +access inside the cluster itself. On modern clusters this means an EKS access entry; +on older clusters it means an entry in the ``aws-auth`` config map. Without that +in-cluster grant the scheduler can read the cluster description but every Kubernetes +API call is rejected. Review Comment: Can this case be avoided by code checks? ########## providers/amazon/docs/executors/eks-executor.rst: ########## @@ -0,0 +1,232 @@ + .. 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. + +.. |executorName| replace:: EKS + +.. _eks_executor: + +================ +AWS EKS Executor +================ + +The EKS executor runs each Airflow task in its own pod on an Amazon EKS cluster. + +This executor extends the Kubernetes executor that ships in the ``cncf.kubernetes`` +provider. Pod scheduling, pod templates, per-task pod overrides, and log handling all +behave the same way they do under +:doc:`apache-airflow-providers-cncf-kubernetes:kubernetes_executor`. The part this executor +adds is authentication. It builds the Kubernetes client for your EKS cluster from an +Airflow AWS connection, so the scheduler does not need a kubeconfig file on disk and +you do not have to refresh cluster credentials yourself. + +Because the executor inherits its pod behavior, everything written about the +Kubernetes executor still applies, including the requirement for a database backend +other than SQLite. + +For a quick start guide please see :ref:`here <eks_setup_guide>`. + +Requirements +------------ + +The executor needs ``apache-airflow-providers-cncf-kubernetes`` version 10.24.0 or +newer, the release that added the ``client_factory`` setting the executor is built on. +Install the ``cncf.kubernetes`` extra of the Amazon provider together with that +version: + +.. code-block:: bash + + pip install 'apache-airflow-providers-amazon[cncf.kubernetes]' \ + 'apache-airflow-providers-cncf-kubernetes>=10.24.0' + +The extra by itself declares a much older floor, because the Amazon provider's other +Kubernetes integrations still work with it. A fresh install resolves to the newest +``cncf.kubernetes`` regardless, so naming the version matters when an older one is +already pinned in your environment. If the installed version is too old, the executor +raises an error when it starts up rather than falling back to a client it cannot +authenticate with. + Review Comment: Lets drop this paragraph. this should be checked with the code. If someone imports the executor without a compatible k8s provider version then just raise exception and explain this in the exception. ########## providers/amazon/docs/executors/eks-executor.rst: ########## @@ -0,0 +1,232 @@ + .. 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. + +.. |executorName| replace:: EKS + +.. _eks_executor: + +================ +AWS EKS Executor +================ + +The EKS executor runs each Airflow task in its own pod on an Amazon EKS cluster. + +This executor extends the Kubernetes executor that ships in the ``cncf.kubernetes`` +provider. Pod scheduling, pod templates, per-task pod overrides, and log handling all +behave the same way they do under +:doc:`apache-airflow-providers-cncf-kubernetes:kubernetes_executor`. The part this executor +adds is authentication. It builds the Kubernetes client for your EKS cluster from an +Airflow AWS connection, so the scheduler does not need a kubeconfig file on disk and +you do not have to refresh cluster credentials yourself. + +Because the executor inherits its pod behavior, everything written about the +Kubernetes executor still applies, including the requirement for a database backend +other than SQLite. + +For a quick start guide please see :ref:`here <eks_setup_guide>`. + +Requirements +------------ + +The executor needs ``apache-airflow-providers-cncf-kubernetes`` version 10.24.0 or +newer, the release that added the ``client_factory`` setting the executor is built on. +Install the ``cncf.kubernetes`` extra of the Amazon provider together with that +version: + +.. code-block:: bash + + pip install 'apache-airflow-providers-amazon[cncf.kubernetes]' \ + 'apache-airflow-providers-cncf-kubernetes>=10.24.0' + +The extra by itself declares a much older floor, because the Amazon provider's other +Kubernetes integrations still work with it. A fresh install resolves to the newest +``cncf.kubernetes`` regardless, so naming the version matters when an older one is +already pinned in your environment. If the installed version is too old, the executor +raises an error when it starts up rather than falling back to a client it cannot +authenticate with. + +The AWS credentials the executor uses must be allowed to call ``eks:DescribeCluster`` +on the cluster, and the IAM principal behind those credentials must be granted +access inside the cluster itself. On modern clusters this means an EKS access entry; +on older clusters it means an entry in the ``aws-auth`` config map. Without that +in-cluster grant the scheduler can read the cluster description but every Kubernetes +API call is rejected. + +How authentication works +------------------------ + +When the executor starts, it plugs its own Kubernetes client into the Kubernetes +executor. Every process that needs a client, including the pod watcher that Airflow runs +as a separate process, builds its own from the ``[aws_eks_executor]`` settings. + +To build a client, the executor describes the cluster to find its API endpoint and +certificate authority, then mints an authentication token for it. An EKS token is a presigned +STS URL that stays valid for roughly fifteen minutes, while the scheduler holds a +single client for as long as it runs. To keep the token current, the executor +registers a refresh hook that the Kubernetes client calls on every authenticated +request. Minting a token is a local signing operation with no network call, so +refreshing that often costs very little. Rotation of the underlying AWS credentials +is handled by botocore in the usual way. Review Comment: what is usual way? This is ambiguous. lets avoid describing it and just replace this with a link to the relevant botocore docs where they explain how they do it? ########## providers/amazon/docs/executors/eks-executor.rst: ########## @@ -0,0 +1,232 @@ + .. 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. + +.. |executorName| replace:: EKS + +.. _eks_executor: + +================ +AWS EKS Executor +================ + +The EKS executor runs each Airflow task in its own pod on an Amazon EKS cluster. + +This executor extends the Kubernetes executor that ships in the ``cncf.kubernetes`` +provider. Pod scheduling, pod templates, per-task pod overrides, and log handling all +behave the same way they do under +:doc:`apache-airflow-providers-cncf-kubernetes:kubernetes_executor`. The part this executor +adds is authentication. It builds the Kubernetes client for your EKS cluster from an +Airflow AWS connection, so the scheduler does not need a kubeconfig file on disk and +you do not have to refresh cluster credentials yourself. + +Because the executor inherits its pod behavior, everything written about the +Kubernetes executor still applies, including the requirement for a database backend +other than SQLite. + +For a quick start guide please see :ref:`here <eks_setup_guide>`. + +Requirements +------------ + +The executor needs ``apache-airflow-providers-cncf-kubernetes`` version 10.24.0 or +newer, the release that added the ``client_factory`` setting the executor is built on. +Install the ``cncf.kubernetes`` extra of the Amazon provider together with that +version: + +.. code-block:: bash + + pip install 'apache-airflow-providers-amazon[cncf.kubernetes]' \ + 'apache-airflow-providers-cncf-kubernetes>=10.24.0' + +The extra by itself declares a much older floor, because the Amazon provider's other +Kubernetes integrations still work with it. A fresh install resolves to the newest +``cncf.kubernetes`` regardless, so naming the version matters when an older one is +already pinned in your environment. If the installed version is too old, the executor +raises an error when it starts up rather than falling back to a client it cannot +authenticate with. + +The AWS credentials the executor uses must be allowed to call ``eks:DescribeCluster`` +on the cluster, and the IAM principal behind those credentials must be granted +access inside the cluster itself. On modern clusters this means an EKS access entry; +on older clusters it means an entry in the ``aws-auth`` config map. Without that +in-cluster grant the scheduler can read the cluster description but every Kubernetes +API call is rejected. + +How authentication works +------------------------ + +When the executor starts, it plugs its own Kubernetes client into the Kubernetes +executor. Every process that needs a client, including the pod watcher that Airflow runs +as a separate process, builds its own from the ``[aws_eks_executor]`` settings. + +To build a client, the executor describes the cluster to find its API endpoint and +certificate authority, then mints an authentication token for it. An EKS token is a presigned +STS URL that stays valid for roughly fifteen minutes, while the scheduler holds a +single client for as long as it runs. To keep the token current, the executor +registers a refresh hook that the Kubernetes client calls on every authenticated +request. Minting a token is a local signing operation with no network call, so +refreshing that often costs very little. Rotation of the underlying AWS credentials +is handled by botocore in the usual way. + +Before it accepts any tasks, the executor checks that the cluster is ``ACTIVE`` (or +``UPDATING``) and, unless ``check_health_on_startup`` is turned off, that it is allowed +to list pods in its namespace. It refuses to start if either check fails. Review Comment: why updating is valid? wha is check_health_on_startup? lets avoid referencing a setting as hard string. can you link to the setting section? ########## providers/amazon/docs/executors/eks-executor.rst: ########## @@ -0,0 +1,232 @@ + .. 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. + +.. |executorName| replace:: EKS + +.. _eks_executor: + +================ +AWS EKS Executor +================ + +The EKS executor runs each Airflow task in its own pod on an Amazon EKS cluster. + +This executor extends the Kubernetes executor that ships in the ``cncf.kubernetes`` +provider. Pod scheduling, pod templates, per-task pod overrides, and log handling all +behave the same way they do under +:doc:`apache-airflow-providers-cncf-kubernetes:kubernetes_executor`. The part this executor +adds is authentication. It builds the Kubernetes client for your EKS cluster from an +Airflow AWS connection, so the scheduler does not need a kubeconfig file on disk and +you do not have to refresh cluster credentials yourself. + +Because the executor inherits its pod behavior, everything written about the +Kubernetes executor still applies, including the requirement for a database backend +other than SQLite. + +For a quick start guide please see :ref:`here <eks_setup_guide>`. + +Requirements +------------ + +The executor needs ``apache-airflow-providers-cncf-kubernetes`` version 10.24.0 or +newer, the release that added the ``client_factory`` setting the executor is built on. +Install the ``cncf.kubernetes`` extra of the Amazon provider together with that +version: + +.. code-block:: bash + + pip install 'apache-airflow-providers-amazon[cncf.kubernetes]' \ + 'apache-airflow-providers-cncf-kubernetes>=10.24.0' + +The extra by itself declares a much older floor, because the Amazon provider's other +Kubernetes integrations still work with it. A fresh install resolves to the newest +``cncf.kubernetes`` regardless, so naming the version matters when an older one is +already pinned in your environment. If the installed version is too old, the executor +raises an error when it starts up rather than falling back to a client it cannot +authenticate with. + +The AWS credentials the executor uses must be allowed to call ``eks:DescribeCluster`` +on the cluster, and the IAM principal behind those credentials must be granted +access inside the cluster itself. On modern clusters this means an EKS access entry; +on older clusters it means an entry in the ``aws-auth`` config map. Without that +in-cluster grant the scheduler can read the cluster description but every Kubernetes +API call is rejected. + +How authentication works +------------------------ + +When the executor starts, it plugs its own Kubernetes client into the Kubernetes +executor. Every process that needs a client, including the pod watcher that Airflow runs +as a separate process, builds its own from the ``[aws_eks_executor]`` settings. + +To build a client, the executor describes the cluster to find its API endpoint and +certificate authority, then mints an authentication token for it. An EKS token is a presigned +STS URL that stays valid for roughly fifteen minutes, while the scheduler holds a +single client for as long as it runs. To keep the token current, the executor +registers a refresh hook that the Kubernetes client calls on every authenticated +request. Minting a token is a local signing operation with no network call, so +refreshing that often costs very little. Rotation of the underlying AWS credentials +is handled by botocore in the usual way. + +Before it accepts any tasks, the executor checks that the cluster is ``ACTIVE`` (or +``UPDATING``) and, unless ``check_health_on_startup`` is turned off, that it is allowed +to list pods in its namespace. It refuses to start if either check fails. + +.. _eks_config_options: + +Config Options +-------------- + +The executor reads its own settings from an ``aws_eks_executor`` section in +``airflow.cfg``. You can also set any of them with an environment variable using the +``AIRFLOW__AWS_EKS_EXECUTOR__<OPTION_NAME>`` form, for example +``AIRFLOW__AWS_EKS_EXECUTOR__CLUSTER_NAME=airflow-eks-cluster``. For more information on +how to set these options, see `Setting Configuration Options +<https://airflow.apache.org/docs/apache-airflow/stable/howto/set-config.html>`__. Review Comment: internal links should not be with http. if doc changes location we won't know it. ########## providers/amazon/src/airflow/providers/amazon/aws/executors/eks/eks_executor.py: ########## @@ -0,0 +1,145 @@ +# 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. +"""AwsEksExecutor: run Airflow tasks as pods on an Amazon EKS cluster.""" + +from __future__ import annotations + +import os +from typing import TYPE_CHECKING + +from airflow.providers.common.compat.sdk import conf + +try: + from kubernetes.client.rest import ApiException + + from airflow.providers.cncf.kubernetes import __version__ as cncf_kubernetes_version + from airflow.providers.cncf.kubernetes.executors.kubernetes_executor import KubernetesExecutor + from airflow.providers.cncf.kubernetes.get_provider_info import ( + get_provider_info as get_cncf_kubernetes_provider_info, + ) +except ImportError as e: + raise ImportError( + "AwsEksExecutor requires the cncf.kubernetes provider; install it with " + "pip install 'apache-airflow-providers-amazon[cncf.kubernetes]'" Review Comment: probably best also to mention min version of supported cncf in the note? ########## providers/amazon/src/airflow/providers/amazon/aws/executors/eks/_client_factory.py: ########## @@ -0,0 +1,130 @@ +# 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. +"""Kubernetes client factories for the AwsEksExecutor.""" + +# Internal to AwsEksExecutor, not a public API. The executor writes the import paths of +# _get_eks_kube_client and _get_eks_async_kube_client into [kubernetes_executor] client_factory +# and async_client_factory, and cncf.kubernetes re-resolves them by that path in every process +# that builds a client, including the spawned pod watcher. Keep eks_executor's paths in step. + +from __future__ import annotations + +import os +import tempfile +from base64 import b64decode +from typing import TYPE_CHECKING + +from airflow.providers.amazon.aws.executors.eks.utils import ( + CONFIG_DEFAULTS, + CONFIG_GROUP_NAME, + AllEksConfigKeys, +) +from airflow.providers.amazon.aws.hooks.eks import EksHook +from airflow.providers.amazon.aws.hooks.sts import StsHook +from airflow.providers.amazon.aws.utils.eks_get_token import fetch_access_token_for_cluster +from airflow.providers.common.compat.sdk import conf + +if TYPE_CHECKING: + from kubernetes import client + from kubernetes_asyncio import client as async_client + +# UPDATING still serves the Kubernetes API, so only creating, deleting and failed clusters are refused. +_USABLE_CLUSTER_STATUSES = ("ACTIVE", "UPDATING") + + +def _get_eks_kube_client() -> client.CoreV1Api: + """Build a Kubernetes client for the configured EKS cluster, in memory and without a kubeconfig.""" + from kubernetes import client + + configuration = client.Configuration() + _configure_eks_auth(configuration) + return client.CoreV1Api(client.ApiClient(configuration=configuration)) + + +def _get_eks_async_kube_client() -> async_client.CoreV1Api: + """Build the asynchronous Kubernetes client used when ``async_pod_creation`` is enabled.""" + from kubernetes_asyncio import client as async_client + + configuration = async_client.Configuration() + _configure_eks_auth(configuration) + return async_client.CoreV1Api(async_client.ApiClient(configuration=configuration)) + + +def _configure_eks_auth(configuration: client.Configuration | async_client.Configuration) -> None: + cluster_name = conf.get(CONFIG_GROUP_NAME, AllEksConfigKeys.CLUSTER_NAME, fallback=None) + if not cluster_name: + raise ValueError(f"[{CONFIG_GROUP_NAME}] cluster_name is required to build an EKS client") + region_name = conf.get(CONFIG_GROUP_NAME, AllEksConfigKeys.REGION_NAME, fallback=None) + conn_id = conf.get( + CONFIG_GROUP_NAME, + AllEksConfigKeys.AWS_CONN_ID, + fallback=CONFIG_DEFAULTS[AllEksConfigKeys.AWS_CONN_ID], + ) + + eks_hook = EksHook(aws_conn_id=conn_id, region_name=region_name) + cluster = eks_hook.conn.describe_cluster(name=cluster_name)["cluster"] + if cluster["status"] not in _USABLE_CLUSTER_STATUSES: + raise ValueError( + f"EKS cluster {cluster_name} is {cluster['status']}; the executor needs it to be ACTIVE" + ) + session = eks_hook.get_session() + + # EKS only accepts tokens presigned against the regional STS endpoint; some regions + # otherwise default to the global one. Same dance as EksHook.generate_config_file. + os.environ["AWS_STS_REGIONAL_ENDPOINTS"] = "regional" + try: + sts_endpoint = StsHook( + aws_conn_id=conn_id, region_name=session.region_name + ).conn_client_meta.endpoint_url + finally: + del os.environ["AWS_STS_REGIONAL_ENDPOINTS"] + sts_url = f"{sts_endpoint}/?Action=GetCallerIdentity&Version=2011-06-15" + + configuration.host = cluster["endpoint"] + configuration.ssl_ca_cert = _write_cluster_ca_file(cluster["certificateAuthority"]["data"]) + # Key the bearer auth under both identifiers. kubernetes-client >= 36 looks the token up under + # "BearerToken" and only aliases the api_key (not api_key_prefix) back to the legacy + # "authorization" slot (see kubernetes-client/python#2595); older clients use "authorization". + # Setting the prefix only under "authorization" makes the newer client emit the raw token with + # no "Bearer " prefix, which the API server rejects with 401, so set both slots. + for identifier in ("BearerToken", "authorization"): + configuration.api_key_prefix[identifier] = "Bearer" Review Comment: This is too much wording. We can avoid this all together by checking the underlying library version. If >=36 do X else do Y that way once we bump min library version would also be easier to adjust the code. ########## providers/amazon/src/airflow/providers/amazon/aws/executors/eks/eks_executor.py: ########## @@ -0,0 +1,145 @@ +# 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. +"""AwsEksExecutor: run Airflow tasks as pods on an Amazon EKS cluster.""" + +from __future__ import annotations + +import os +from typing import TYPE_CHECKING + +from airflow.providers.common.compat.sdk import conf + +try: + from kubernetes.client.rest import ApiException + + from airflow.providers.cncf.kubernetes import __version__ as cncf_kubernetes_version + from airflow.providers.cncf.kubernetes.executors.kubernetes_executor import KubernetesExecutor + from airflow.providers.cncf.kubernetes.get_provider_info import ( + get_provider_info as get_cncf_kubernetes_provider_info, + ) +except ImportError as e: + raise ImportError( + "AwsEksExecutor requires the cncf.kubernetes provider; install it with " + "pip install 'apache-airflow-providers-amazon[cncf.kubernetes]'" + ) from e + +from airflow.providers.amazon.aws.executors.eks.utils import ( + CONFIG_DEFAULTS, + CONFIG_GROUP_NAME, + AllEksConfigKeys, +) + +# Import paths, because cncf.kubernetes re-resolves the factories in each process (see _client_factory). +_FACTORY_MODULE = "airflow.providers.amazon.aws.executors.eks._client_factory" +_CLIENT_FACTORY_PATH = f"{_FACTORY_MODULE}._get_eks_kube_client" +_ASYNC_CLIENT_FACTORY_PATH = f"{_FACTORY_MODULE}._get_eks_async_kube_client" +# Only used in the error message below. The check looks for the client_factory option instead of +# comparing versions, because an unreleased source tree still reports the previous release. +MIN_CNCF_KUBERNETES_VERSION = "10.24.0" + + +class AwsEksExecutor(KubernetesExecutor): + """ + A KubernetesExecutor that authenticates against an Amazon EKS cluster. + + Builds the Kubernetes client from ``[aws_eks_executor]`` configuration and keeps the + short-lived EKS token fresh in-process. All pod-level behaviour comes unchanged from the + KubernetesExecutor and its ``[kubernetes_executor]`` configuration. + """ + + # The client factories read the un-prefixed [aws_eks_executor] section, so every team + # would land on the same cluster. + supports_multi_team: bool = False + + def __init__(self, *args, **kwargs): + self._validate_eks_config() + self._require_client_factory_support() + self._ensure_client_factory() + super().__init__(*args, **kwargs) + + def start(self) -> None: + """Call this when the Executor is run for the first time by the scheduler.""" + super().start() + check_health = conf.getboolean( + CONFIG_GROUP_NAME, + AllEksConfigKeys.CHECK_HEALTH_ON_STARTUP, + fallback=CONFIG_DEFAULTS[AllEksConfigKeys.CHECK_HEALTH_ON_STARTUP], + ) + + if not check_health: + return + + self.log.info("Starting EKS Executor and determining health...") + try: + self.check_health() + except RuntimeError: + self.log.error("Stopping the Airflow Scheduler from starting until the issue is resolved.") + raise + + def check_health(self) -> None: + """ + Make a test Kubernetes API call to check the health of the EKS Executor. + + Building the client only proves the cluster exists; without this, a missing EKS access + entry or RBAC binding only shows up later as a watcher error loop. + """ + if TYPE_CHECKING: + assert self.kube_client + namespace = self.kube_config.kube_namespace + try: + self.kube_client.list_namespaced_pod(namespace, limit=1) + except ApiException as e: + raise RuntimeError( + f"EKS Executor health check has failed because: cannot list pods in namespace {namespace} " + f"({e.status} {e.reason}). Check the EKS access entry and RBAC for the executor's IAM role." + ) from e + self.log.info("EKS Executor health check has succeeded.") + + @staticmethod + def _require_client_factory_support() -> None: + # A cncf.kubernetes without the seam imports fine but never reads client_factory, so it + # would build a default client from kubeconfig and fail later as an authentication error. + provider_config = get_cncf_kubernetes_provider_info().get("config", {}) + if "client_factory" not in provider_config.get("kubernetes_executor", {}).get("options", {}): + raise ImportError( + "AwsEksExecutor requires apache-airflow-providers-cncf-kubernetes>=" + f"{MIN_CNCF_KUBERNETES_VERSION}, which honors the [kubernetes_executor] " + f"client_factory setting. The installed {cncf_kubernetes_version} ignores it." + ) Review Comment: I prefer that all checks against supported versions be done at the import level If we don't have a good version we should never reach to this function. That way functions should not worry about supported versions checks. ########## providers/amazon/docs/executors/eks-executor.rst: ########## Review Comment: I think we need also an edit to the AWS executors docs to answer the question: if I am AWS user who want to understand my options how do I choose between Eks vs Lambda vs ECS? ########## providers/amazon/src/airflow/providers/amazon/aws/executors/eks/_client_factory.py: ########## @@ -0,0 +1,130 @@ +# 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. +"""Kubernetes client factories for the AwsEksExecutor.""" + +# Internal to AwsEksExecutor, not a public API. The executor writes the import paths of +# _get_eks_kube_client and _get_eks_async_kube_client into [kubernetes_executor] client_factory +# and async_client_factory, and cncf.kubernetes re-resolves them by that path in every process +# that builds a client, including the spawned pod watcher. Keep eks_executor's paths in step. + +from __future__ import annotations + +import os +import tempfile +from base64 import b64decode +from typing import TYPE_CHECKING + +from airflow.providers.amazon.aws.executors.eks.utils import ( + CONFIG_DEFAULTS, + CONFIG_GROUP_NAME, + AllEksConfigKeys, +) +from airflow.providers.amazon.aws.hooks.eks import EksHook +from airflow.providers.amazon.aws.hooks.sts import StsHook +from airflow.providers.amazon.aws.utils.eks_get_token import fetch_access_token_for_cluster +from airflow.providers.common.compat.sdk import conf + +if TYPE_CHECKING: + from kubernetes import client + from kubernetes_asyncio import client as async_client + +# UPDATING still serves the Kubernetes API, so only creating, deleting and failed clusters are refused. +_USABLE_CLUSTER_STATUSES = ("ACTIVE", "UPDATING") + + +def _get_eks_kube_client() -> client.CoreV1Api: + """Build a Kubernetes client for the configured EKS cluster, in memory and without a kubeconfig.""" + from kubernetes import client + + configuration = client.Configuration() + _configure_eks_auth(configuration) + return client.CoreV1Api(client.ApiClient(configuration=configuration)) + + +def _get_eks_async_kube_client() -> async_client.CoreV1Api: + """Build the asynchronous Kubernetes client used when ``async_pod_creation`` is enabled.""" + from kubernetes_asyncio import client as async_client + + configuration = async_client.Configuration() + _configure_eks_auth(configuration) + return async_client.CoreV1Api(async_client.ApiClient(configuration=configuration)) + + +def _configure_eks_auth(configuration: client.Configuration | async_client.Configuration) -> None: + cluster_name = conf.get(CONFIG_GROUP_NAME, AllEksConfigKeys.CLUSTER_NAME, fallback=None) + if not cluster_name: + raise ValueError(f"[{CONFIG_GROUP_NAME}] cluster_name is required to build an EKS client") + region_name = conf.get(CONFIG_GROUP_NAME, AllEksConfigKeys.REGION_NAME, fallback=None) + conn_id = conf.get( + CONFIG_GROUP_NAME, + AllEksConfigKeys.AWS_CONN_ID, + fallback=CONFIG_DEFAULTS[AllEksConfigKeys.AWS_CONN_ID], + ) + + eks_hook = EksHook(aws_conn_id=conn_id, region_name=region_name) + cluster = eks_hook.conn.describe_cluster(name=cluster_name)["cluster"] + if cluster["status"] not in _USABLE_CLUSTER_STATUSES: + raise ValueError( + f"EKS cluster {cluster_name} is {cluster['status']}; the executor needs it to be ACTIVE" + ) Review Comment: Docs mention that cluster also accept tasks in `UPDATING` state but here we check for active. I assume doc entry is wrong? ########## providers/amazon/src/airflow/providers/amazon/aws/executors/eks/_client_factory.py: ########## @@ -0,0 +1,130 @@ +# 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. +"""Kubernetes client factories for the AwsEksExecutor.""" + +# Internal to AwsEksExecutor, not a public API. The executor writes the import paths of +# _get_eks_kube_client and _get_eks_async_kube_client into [kubernetes_executor] client_factory +# and async_client_factory, and cncf.kubernetes re-resolves them by that path in every process +# that builds a client, including the spawned pod watcher. Keep eks_executor's paths in step. + +from __future__ import annotations + +import os +import tempfile +from base64 import b64decode +from typing import TYPE_CHECKING + +from airflow.providers.amazon.aws.executors.eks.utils import ( + CONFIG_DEFAULTS, + CONFIG_GROUP_NAME, + AllEksConfigKeys, +) +from airflow.providers.amazon.aws.hooks.eks import EksHook +from airflow.providers.amazon.aws.hooks.sts import StsHook +from airflow.providers.amazon.aws.utils.eks_get_token import fetch_access_token_for_cluster +from airflow.providers.common.compat.sdk import conf + +if TYPE_CHECKING: + from kubernetes import client + from kubernetes_asyncio import client as async_client + +# UPDATING still serves the Kubernetes API, so only creating, deleting and failed clusters are refused. +_USABLE_CLUSTER_STATUSES = ("ACTIVE", "UPDATING") + + +def _get_eks_kube_client() -> client.CoreV1Api: + """Build a Kubernetes client for the configured EKS cluster, in memory and without a kubeconfig.""" + from kubernetes import client + + configuration = client.Configuration() + _configure_eks_auth(configuration) + return client.CoreV1Api(client.ApiClient(configuration=configuration)) + + +def _get_eks_async_kube_client() -> async_client.CoreV1Api: + """Build the asynchronous Kubernetes client used when ``async_pod_creation`` is enabled.""" + from kubernetes_asyncio import client as async_client + + configuration = async_client.Configuration() + _configure_eks_auth(configuration) + return async_client.CoreV1Api(async_client.ApiClient(configuration=configuration)) + + +def _configure_eks_auth(configuration: client.Configuration | async_client.Configuration) -> None: + cluster_name = conf.get(CONFIG_GROUP_NAME, AllEksConfigKeys.CLUSTER_NAME, fallback=None) + if not cluster_name: + raise ValueError(f"[{CONFIG_GROUP_NAME}] cluster_name is required to build an EKS client") + region_name = conf.get(CONFIG_GROUP_NAME, AllEksConfigKeys.REGION_NAME, fallback=None) + conn_id = conf.get( + CONFIG_GROUP_NAME, + AllEksConfigKeys.AWS_CONN_ID, + fallback=CONFIG_DEFAULTS[AllEksConfigKeys.AWS_CONN_ID], + ) + + eks_hook = EksHook(aws_conn_id=conn_id, region_name=region_name) + cluster = eks_hook.conn.describe_cluster(name=cluster_name)["cluster"] + if cluster["status"] not in _USABLE_CLUSTER_STATUSES: + raise ValueError( + f"EKS cluster {cluster_name} is {cluster['status']}; the executor needs it to be ACTIVE" + ) + session = eks_hook.get_session() + + # EKS only accepts tokens presigned against the regional STS endpoint; some regions + # otherwise default to the global one. Same dance as EksHook.generate_config_file. + os.environ["AWS_STS_REGIONAL_ENDPOINTS"] = "regional" + try: + sts_endpoint = StsHook( + aws_conn_id=conn_id, region_name=session.region_name + ).conn_client_meta.endpoint_url + finally: + del os.environ["AWS_STS_REGIONAL_ENDPOINTS"] + sts_url = f"{sts_endpoint}/?Action=GetCallerIdentity&Version=2011-06-15" + + configuration.host = cluster["endpoint"] + configuration.ssl_ca_cert = _write_cluster_ca_file(cluster["certificateAuthority"]["data"]) + # Key the bearer auth under both identifiers. kubernetes-client >= 36 looks the token up under + # "BearerToken" and only aliases the api_key (not api_key_prefix) back to the legacy + # "authorization" slot (see kubernetes-client/python#2595); older clients use "authorization". + # Setting the prefix only under "authorization" makes the newer client emit the raw token with + # no "Bearer " prefix, which the API server rejects with 401, so set both slots. + for identifier in ("BearerToken", "authorization"): + configuration.api_key_prefix[identifier] = "Bearer" + + # The EKS token is a presigned STS URL valid for ~15 minutes, but the scheduler holds + # one client for its whole lifetime. refresh_api_key_hook runs on every authenticated + # request (Configuration.get_api_key_with_prefix), and minting is a local SigV4 signing + # with no network round-trip, so re-minting per request keeps the token fresh at + # negligible cost. Botocore refreshes the session's own credentials when they rotate. + def refresh_api_key(config: client.Configuration | async_client.Configuration) -> None: + token = fetch_access_token_for_cluster( + cluster_name, sts_url, region_name=session.region_name, session=session + ) Review Comment: I am not sure I understand this one. Shouldn't the refresh request be in a side car container? ########## providers/amazon/docs/executors/eks-executor.rst: ########## @@ -0,0 +1,232 @@ + .. 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. + +.. |executorName| replace:: EKS + +.. _eks_executor: + +================ +AWS EKS Executor +================ + +The EKS executor runs each Airflow task in its own pod on an Amazon EKS cluster. + +This executor extends the Kubernetes executor that ships in the ``cncf.kubernetes`` +provider. Pod scheduling, pod templates, per-task pod overrides, and log handling all +behave the same way they do under +:doc:`apache-airflow-providers-cncf-kubernetes:kubernetes_executor`. The part this executor +adds is authentication. It builds the Kubernetes client for your EKS cluster from an +Airflow AWS connection, so the scheduler does not need a kubeconfig file on disk and +you do not have to refresh cluster credentials yourself. + +Because the executor inherits its pod behavior, everything written about the +Kubernetes executor still applies, including the requirement for a database backend +other than SQLite. + +For a quick start guide please see :ref:`here <eks_setup_guide>`. + +Requirements +------------ + +The executor needs ``apache-airflow-providers-cncf-kubernetes`` version 10.24.0 or +newer, the release that added the ``client_factory`` setting the executor is built on. +Install the ``cncf.kubernetes`` extra of the Amazon provider together with that +version: + +.. code-block:: bash + + pip install 'apache-airflow-providers-amazon[cncf.kubernetes]' \ + 'apache-airflow-providers-cncf-kubernetes>=10.24.0' + +The extra by itself declares a much older floor, because the Amazon provider's other +Kubernetes integrations still work with it. A fresh install resolves to the newest +``cncf.kubernetes`` regardless, so naming the version matters when an older one is +already pinned in your environment. If the installed version is too old, the executor +raises an error when it starts up rather than falling back to a client it cannot +authenticate with. + +The AWS credentials the executor uses must be allowed to call ``eks:DescribeCluster`` +on the cluster, and the IAM principal behind those credentials must be granted +access inside the cluster itself. On modern clusters this means an EKS access entry; +on older clusters it means an entry in the ``aws-auth`` config map. Without that +in-cluster grant the scheduler can read the cluster description but every Kubernetes +API call is rejected. + +How authentication works +------------------------ + +When the executor starts, it plugs its own Kubernetes client into the Kubernetes +executor. Every process that needs a client, including the pod watcher that Airflow runs +as a separate process, builds its own from the ``[aws_eks_executor]`` settings. + +To build a client, the executor describes the cluster to find its API endpoint and +certificate authority, then mints an authentication token for it. An EKS token is a presigned +STS URL that stays valid for roughly fifteen minutes, while the scheduler holds a +single client for as long as it runs. To keep the token current, the executor +registers a refresh hook that the Kubernetes client calls on every authenticated +request. Minting a token is a local signing operation with no network call, so +refreshing that often costs very little. Rotation of the underlying AWS credentials +is handled by botocore in the usual way. + +Before it accepts any tasks, the executor checks that the cluster is ``ACTIVE`` (or +``UPDATING``) and, unless ``check_health_on_startup`` is turned off, that it is allowed +to list pods in its namespace. It refuses to start if either check fails. + +.. _eks_config_options: + +Config Options +-------------- + +The executor reads its own settings from an ``aws_eks_executor`` section in +``airflow.cfg``. You can also set any of them with an environment variable using the +``AIRFLOW__AWS_EKS_EXECUTOR__<OPTION_NAME>`` form, for example +``AIRFLOW__AWS_EKS_EXECUTOR__CLUSTER_NAME=airflow-eks-cluster``. For more information on +how to set these options, see `Setting Configuration Options +<https://airflow.apache.org/docs/apache-airflow/stable/howto/set-config.html>`__. + +Required config options: +~~~~~~~~~~~~~~~~~~~~~~~~ + +- CLUSTER_NAME - The name of the Amazon EKS cluster that tasks run on. The + executor refuses to start if this is unset. Required. + +Optional config options: +~~~~~~~~~~~~~~~~~~~~~~~~ + +- CONN_ID - The Airflow connection (i.e. credentials) used by the EKS + executor to make API calls to Amazon EKS. Defaults to ``aws_default``. +- REGION_NAME - The AWS Region the cluster is in. When this is left empty, the + region comes from the standard boto3 resolution order. +- CHECK_HEALTH_ON_STARTUP - Whether to check on startup that the executor can + list pods in its namespace. Defaults to ``True``. + +Pod-level configuration +~~~~~~~~~~~~~~~~~~~~~~~ + +Settings that describe the worker pods themselves stay in the +``[kubernetes_executor]`` section, exactly as they are for the Kubernetes executor. +This includes ``namespace``, ``pod_template_file``, ``worker_container_repository``, +``worker_container_tag``, and ``delete_worker_pods``. See +:doc:`apache-airflow-providers-cncf-kubernetes:kubernetes_executor` for the full list +and for how pod templates and ``pod_override`` work. + +Leave ``client_factory`` and ``async_client_factory`` in that section unset. The +executor sets both itself, and raises an error at startup if either is set to something +else, so that a stale or conflicting setting cannot silently send tasks to the wrong +cluster. + +The worker pod template +~~~~~~~~~~~~~~~~~~~~~~~ + +Set ``[kubernetes_executor] pod_template_file`` to the path of a YAML file describing +the worker pod. Treat it as required. There is no usable default: when the setting is +empty the Kubernetes executor only logs a warning that the model file does not exist +and then builds a worker pod from an empty template. That pod is missing the container +the executor needs, so the failure arrives later and somewhere else, usually as a +rejected pod or a task that never reports back. + +The file must define a container named ``base`` as the first entry in +``spec.containers``, and that container has to run your Airflow worker image. +The ``pod_template_file`` section of +:doc:`apache-airflow-providers-cncf-kubernetes:kubernetes_executor` covers the full set +of requirements and includes templates for Dags baked into the image, Dags on a volume, +and git-sync. Any of those works here unchanged, since this executor only replaces how +the Kubernetes client is authenticated. + +.. _eks_logging: + +.. include:: general.rst + :start-after: .. BEGIN LOGGING + :end-before: .. END LOGGING + +- Worker pods are deleted once their task finishes (``[kubernetes_executor] + delete_worker_pods``), and their logs go with them, so configure remote + logging to CloudWatch Logs or S3 to keep task logs viewable in the Airflow UI. +- The remote logging configuration must be set on the worker pods as well as on + the scheduler and API server, and the worker pods need an IAM role, for + example through EKS Pod Identity, that can write to the log destination. + +.. _eks_setup_guide: + +Setting up an EKS Executor for Apache Airflow +--------------------------------------------- + +Grant access to the cluster +~~~~~~~~~~~~~~~~~~~~~~~~~~~ + +Create an EKS access entry for the IAM role or user that the scheduler runs as, and +associate a policy that allows it to create and watch pods in the namespace you plan +to use. If you manage cluster access through the ``aws-auth`` config map instead, add +the principal there and map it to a Kubernetes group with the same permissions. + +Configure Airflow +~~~~~~~~~~~~~~~~~ + +Select the executor and name your cluster: + +.. code-block:: ini + + [core] + executor = airflow.providers.amazon.aws.executors.eks.eks_executor.AwsEksExecutor + + [aws_eks_executor] + cluster_name = airflow-eks-cluster + region_name = us-east-1 + conn_id = aws_default + + [kubernetes_executor] + namespace = airflow + pod_template_file = /opt/airflow/pod_templates/worker_template.yaml + worker_container_repository = my-account.dkr.ecr.us-east-1.amazonaws.com/airflow + worker_container_tag = latest + +The worker image must contain Airflow and the Amazon provider, and it needs access +to your Dag files and a network path to the Airflow API server, in the same way any +Kubernetes executor worker does. + +Task logging +~~~~~~~~~~~~ + +Configure remote logging as described in the :ref:`logging <eks_logging>` section, so +that task logs stay viewable in the Airflow UI after the worker pods are deleted. + +Verify the setup +~~~~~~~~~~~~~~~~ + +Start the scheduler and trigger a small Dag. A successful run shows a worker pod +appearing in your chosen namespace and the task finishing in the Airflow UI. If the +scheduler logs an authentication or forbidden error from the Kubernetes API, the +in-cluster access grant for your IAM principal is the first thing to check. + +Multi-team deployments +---------------------- + +The executor cannot be used as a team executor yet, because its settings are read from +the un-prefixed ``[aws_eks_executor]`` section and every team would share one cluster. Review Comment: what is team executor? are you referencing multi team feature? if so please link to the relevant doc section. By `yet` you mean it's planned soon or is it long term plan? if so lets add github issue for users to be able to interact with us on that feature request -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
