View a markdown version of this page

Spark Connect を使用して Amazon EMR on EKS でインタラクティブセッションを実行する - Amazon EMR

翻訳は機械翻訳により提供されています。提供された翻訳内容と英語版の間で齟齬、不一致または矛盾がある場合、英語版が優先します。

Spark Connect を使用して Amazon EMR on EKS でインタラクティブセッションを実行する

Amazon EMR on EKS リリース emr-7.14.0 以降 (emr-spark-8.1.0または 以降) では、VS Code、、PyCharmJupyterノートブックなどのセルフマネージドPySparkクライアントからマネージド Spark Connect エンドポイントに接続できます。Spark Connect は、アプリケーションコードを Spark ドライバープロセスから切り離すクライアント/サーバーアーキテクチャを使用します。Spark オペレーションが Amazon EMR on EKS を介して EKS クラスターで実行される間、ローカル IDE でPySparkコードを開発およびデバッグします。Spark Connect には次の利点があります。

  • VS Code、、PyCharmJupyterノートブックなど、任意のPySparkクライアントから Amazon EMR on EKS に接続する

  • DataFrames が本番環境のデータをリモートで実行している間、IDE でブレークポイントを設定し、PySparkコードをステップスルーする

  • コンピューティング、ネットワーク、セキュリティ設定を完全に制御して独自の EKS クラスターで を実行する

Spark Connect エンドポイントは、Spark Connect サーバーをホストする仮想クラスター上のマネージドエンドポイントです。Spark Connect エンドポイントを作成すると、Amazon EMR on EKS は EKS クラスターに gRPC サーバーを使用して Spark ドライバーをプロビジョニングします。ローカル PySpark クライアントは、gRPC エンドポイントを介して DataFrame および SQL オペレーションをドライバーに送信します。エンドポイントを操作するには、 GetManagedEndpointSessionCredentials API を使用してセッション認証情報を取得します。各エンドポイントは、複数の同時セッションをサポートします。

前提条件

Spark Connect エンドポイントを作成する前に、以下があることを確認してください。

以下の要件は Spark Connect エンドポイントに固有のものです。

  • 少なくとも 1 つのプライベートサブネットを持つ EKS クラスター。gRPC トラフィックをエンドポイントにルーティングするために、Spark Connect は VPC にプライベートサブネットを必要とする内部 Network Load Balancer (NLB) をプロビジョニングします。NLB は内部であり、パブリックインターネットには公開されません。

  • Amazon EMR on EKS が内部 NLB をプロビジョニングできるように、EKS クラスターにインストールされた AWS Load Balancerコントローラー。手順については、AWS Load Balancerコントローラーのインストール」を参照してください。

また、以下の標準 Amazon EMR on EKS リソースも必要です。Amazon EMR on EKS でジョブをすでに実行している場合は、次のような設定になっている可能性があります。

  • Amazon S3 バケットやデータカタログなどのデータソースにアクセスするためのアクセス許可を持つ IAM ジョブ実行ロール。手順については、「ジョブ実行ロールの作成」を参照してください。

  • EKS クラスターがクラスターアクセス管理を使用している場合、Amazon EMR on EKS に必要なアクセスエントリ。手順については、「EKS での Amazon EMR のクラスターアクセスの設定」を参照してください。

必要なアクセス許可

IAM ロールに次のアクセス許可を追加して、Spark Connect エンドポイントを作成して操作します。

{ "Version": "2012-10-17", "Statement": [ { "Sid": "EMRContainersEndpointAccess", "Effect": "Allow", "Action": [ "emr-containers:CreateManagedEndpoint", "emr-containers:DescribeManagedEndpoint", "emr-containers:DeleteManagedEndpoint", "emr-containers:ListManagedEndpoints", "emr-containers:GetManagedEndpointSessionCredentials" ], "Resource": [ "arn:aws:emr-containers:region:account-id:/virtualclusters/virtual-cluster-id", "arn:aws:emr-containers:region:account-id:/virtualclusters/virtual-cluster-id/endpoints/*" ] }, { "Sid": "PassRoleToEKSPodIdentity", "Effect": "Allow", "Action": "iam:PassRole", "Resource": "arn:aws:iam::account-id:role/ExecutionRole", "Condition": { "StringEquals": { "iam:PassedToService": "pods.eks.amazonaws.com" }, "ArnLike": { "iam:AssociatedResourceARN": [ "arn:aws:eks:region:account-id:cluster/eks-cluster-name" ] } } } ] }

セキュリティ設定を作成する

Spark Connect エンドポイントには、仮想クラスターに関連付けられたセキュリティ設定が必要です。セキュリティ設定は、エンドポイントの認証と認可の設定を定義し、Spark Connect インフラストラクチャが実行されるシステム名前空間を指定します。

注記

各セキュリティ設定には、仮想クラスターとの one-to-one の関係があります。複数の仮想クラスターで同じセキュリティ設定を再利用することはできません。

システム名前空間を使用してセキュリティ設定を作成します。

aws emr-containers create-security-configuration \ --name "security-config-name" \ --security-configuration-data '{ "authenticationConfiguration": { "identityCenterConfiguration": { "enableIdentityCenter": false } } }' \ --container-provider '{ "type": "EKS", "id": "eks-cluster-name", "info": { "eksInfo": { "namespace": "system-namespace" } } }'

セキュリティ設定namespaceの は、Spark Connect インフラストラクチャコンポーネントが実行されるシステム名前空間です。これは、仮想クラスターの作成時に指定されたユーザー名前空間とは別のものです。

Spark Connect を有効にして仮想クラスターを作成する

セキュリティ設定に関連付けられた仮想クラスターを作成します。sessionEnabled を に設定trueし、前のステップsecurityConfigurationIdの を指定する必要があります。

aws emr-containers create-virtual-cluster \ --name "virtual-cluster-name" \ --container-provider '{ "type": "EKS", "id": "eks-cluster-name", "info": { "eksInfo": { "namespace": "user-namespace" } } }' \ --security-configuration-id SECURITY_CONFIGURATION_ID \ --session-enabled true
重要

--security-configuration-id — この仮想クラスターを前のステップで作成したセキュリティ設定に関連付けます。これは Spark Connect エンドポイントに必要です。

--session-enabled — 仮想クラスターで Spark Connect エンドポイントのサポートを有効にします。このフラグがないと、この仮想クラスターに Spark Connect エンドポイントを作成することはできません。

Spark Connect エンドポイントを作成する

セキュリティ設定を作成したら、仮想クラスターに Spark Connect マネージドエンドポイントを作成します。

注記

EKS クラスターの最初の Spark Connect エンドポイントは、それ以降のエンドポイントACTIVEよりも長くなります。EKS クラスターの最初のエンドポイントの場合、Amazon EMR on EKS は接続に使用される共有ネットワークコンポーネントをプロビジョニングします。内部 Network Load Balancer (NLB)、VPC インターフェイスエンドポイント (AWS PrivateLink)、およびエンドポイントに gRPC トラフィックをルーティングする 1 回限りの Envoy ルーターデプロイ (spark-connect-router名前空間内) です。これには数分かかる場合があります。これらのコンポーネントは EKS クラスターごとに 1 回のみ作成され、そのクラスターのそれ以降のすべてのエンドポイントで再利用されるため、後続のエンドポイントはより高速に起動します。

aws emr-containers create-managed-endpoint \ --virtual-cluster-id VIRTUAL_CLUSTER_ID \ --name "spark-connect-endpoint" \ --type "SPARK_CONNECT" \ --release-label "emr-7.14.0-latest" \ --execution-role-arn "arn:aws:iam::account-id:role/ExecutionRole" \ --session-idle-timeout-in-minutes 60
注記

デフォルトでは、Spark Connect マネージドエンドポイントは 2 つのエグゼキュターで始まります。設定オーバーライドを指定しない場合、エンドポイントはこのデフォルトを使用します。異なる数のエグゼキュターで を実行するには、次の例に示すように、設定オーバーライドspark.executor.instancesで を設定します。

次の例には、動的割り当てによる Spark 設定の上書きと、Spark ログを Amazon S3 にエクスポートするためのモニタリング設定が含まれています。

aws emr-containers create-managed-endpoint \ --virtual-cluster-id VIRTUAL_CLUSTER_ID \ --name "spark-connect-endpoint" \ --type "SPARK_CONNECT" \ --release-label "emr-7.14.0-latest" \ --execution-role-arn "arn:aws:iam::account-id:role/ExecutionRole" \ --session-idle-timeout-in-minutes 60 \ --configuration-overrides '{ "applicationConfiguration": [ { "classification": "spark-defaults", "properties": { "spark.driver.memory": "4g", "spark.executor.memory": "4g", "spark.executor.cores": "2", "spark.executor.instances": "3", "spark.dynamicAllocation.enabled": "true", "spark.dynamicAllocation.minExecutors": "3", "spark.dynamicAllocation.maxExecutors": "5" } } ], "monitoringConfiguration": { "s3MonitoringConfiguration": { "logUri": "s3://your-bucket/spark-connect-logs/" }, "persistentAppUI": "ENABLED" } }'

エンドポイントのステータスをモニタリングします。

aws emr-containers describe-managed-endpoint \ --virtual-cluster-id VIRTUAL_CLUSTER_ID \ --id ENDPOINT_ID

エンドポイントの状態が になるまで待ってACTIVEから接続します。

Spark Connect エンドポイントに接続する

エンドポイントがアクティブになったら、セッション認証情報を取得し、PySpark クライアントから接続します。

Spark Connect エンドポイントに接続するには
  1. エンドポイントの説明から認証プロキシ URL を取得します。

    aws emr-containers describe-managed-endpoint \ --virtual-cluster-id VIRTUAL_CLUSTER_ID \ --id ENDPOINT_ID

    エンドポイントが の場合、レスポンスには authProxyUrlフィールドが含まれますACTIVE。

  2. エンドポイントのセッショントークンを取得します。

    aws emr-containers get-managed-endpoint-session-credentials \ --virtual-cluster-identifier VIRTUAL_CLUSTER_ID \ --endpoint-identifier ENDPOINT_ID \ --execution-role-arn "arn:aws:iam::account-id:role/ExecutionRole" \ --credential-type "TOKEN"

    レスポンスにはセッショントークンが含まれます。

    { "id": "SESSION_ID", "credentials": { "token": "SESSION_TOKEN" }, "endpointCredentials": { "token": "ENDPOINT_CREDENTIALS_TOKEN" }, "expiresAt": "EXPIRY_TIME" }

    認証プロキシ URL に接続するときは、 credentials.token値を x-aws-proxy-authパラメータとして使用します。

  3. エンドポイントの Spark バージョンに一致する PySpark クライアントをインストールします ( の場合は Spark 3.5.8emr-7.14.0、 の場合は Spark 4.1.1emr-spark-8.1.0)。

    # For emr-7.14.0 pip install pyspark[connect]==3.5.8 # For emr-spark-8.1.0 pip install pyspark[connect]==4.1.1 pip install boto3
  4. 認証プロキシ URL とセッショントークンを使用して PySpark から接続します。

    from pyspark.sql import SparkSession # Use the authProxyUrl from DescribeManagedEndpoint # and the credentials.token from GetManagedEndpointSessionCredentials connect_url = f"{auth_proxy_url}/;use_ssl=true;x-aws-proxy-auth={token}" spark = SparkSession.builder.remote(connect_url).getOrCreate() print(f"Connected. Spark version: {spark.version}") # Run queries spark.sql("SELECT 1+1 AS result").show() # When finished, disconnect the client spark.stop()

次の Python スクリプトは、認証プロキシ URL とセッショントークンを取得し、Spark Connect エンドポイントに接続します。

import boto3 from pyspark.sql import SparkSession from pyspark.sql.functions import col REGION = 'REGION' VIRTUAL_CLUSTER_ID = 'VIRTUAL_CLUSTER_ID' ENDPOINT_ID = 'ENDPOINT_ID' EXECUTION_ROLE = 'arn:aws:iam::account-id:role/ExecutionRole' client = boto3.client('emr-containers', region_name=REGION) # Get the auth proxy URL from DescribeManagedEndpoint endpoint_response = client.describe_managed_endpoint( virtualClusterId=VIRTUAL_CLUSTER_ID, id=ENDPOINT_ID ) auth_proxy_url = endpoint_response['endpoint']['authProxyUrl'] # Get session token creds_response = client.get_managed_endpoint_session_credentials( virtualClusterIdentifier=VIRTUAL_CLUSTER_ID, endpointIdentifier=ENDPOINT_ID, executionRoleArn=EXECUTION_ROLE, credentialType='TOKEN' ) token = creds_response['credentials']['token'] # Connect via Spark Connect using auth proxy URL and credentials token connect_url = f"{auth_proxy_url}/;use_ssl=true;x-aws-proxy-auth={token}" spark = SparkSession.builder.remote(connect_url).getOrCreate() print(f"Connected. Spark version: {spark.version}") # Run DataFrame operations df = spark.range(100).withColumn("squared", col("id") * col("id")) df.show(10) print(f"Count: {df.count()}") spark.stop()

考慮事項と制限事項

Amazon EMR on EKS で Spark Connect を使用してインタラクティブワークロードを実行する場合は、次の点を考慮してください。

  • Spark Connect は、Amazon EMR on EKS リリース emr-7.14.0 以降、または emr-spark-8.1.0以降でサポートされています。

  • Spark Connect は、 で DataFrame および SQL APIs をサポートしていますPySpark。Spark Connect は RDD ベースの APIsをサポートしていません。

  • セッショントークンには時間制限があります。トークンの有効期限が切れると、gRPC 呼び出しは認証エラーで失敗します。GetManagedEndpointSessionCredentials を呼び出して新しいトークンを取得し、更新されたトークンSparkSessionで新しい を作成します。

  • 各セキュリティ設定には、仮想クラスターとの one-to-one の関係があります。セキュリティ設定を複数の仮想クラスター間で共有することはできません。

  • セキュリティ設定を削除する前に、セキュリティ設定を使用してすべてのエンドポイントを削除する必要があります。

  • ローカルにインストールされる PySpark バージョンは、エンドポイントの Apache Spark バージョンと一致する必要があります ( の場合は Spark 3.5.8emr-7.14.0、 の場合は Spark 4.1.1emr-spark-8.1.0)。バージョンが一致しない場合、接続エラーまたは予期しない動作が発生します。

  • Spark Connect エンドポイントタイプは ですSPARK_CONNECT。これは Livy インタラクティブエンドポイント (タイプ ) とは異なりますJUPYTER_ENTERPRISE_GATEWAY。

  • sessionIdleTimeoutInMinutes パラメータは、自動終了までにアイドルセッションが持続する時間を制御します。デフォルトは、60 分です。

  • Spark Connect エンドポイントは、信頼できる ID の伝播をサポートしていません。

  • Spark Connect エンドポイントは、まだ Lake Formation のきめ細かなアクセスコントロール (FGAC) をサポートしていません。アクセスコントロールを適用するには、エンドポイントに関連付けられた IAM 実行ロールを使用します。

  • Spark Connect エンドポイントは、Network Load Balancer (NLB) を使用して gRPC トラフィックをルーティングします。NLB は、最初の Spark Connect エンドポイントが作成されたときに作成され、最後のセッション対応仮想クラスターが削除されたときにのみ削除されます。Amazon EMR on EKS は、gRPC トラフィックをエンドポイントにルーティングする単一の Envoy ルーター (spark-connect-router名前空間では EKS クラスターごとに 1 つ) も実行します。最初の Spark Connect エンドポイントで作成され、最後のセッション対応仮想クラスターが削除されると終了します。セッション中に Spark ドライバーとエグゼキュターによって消費される EKS コンピューティングリソースに加えて、NLB コストと Envoy ルーターによって消費される EKS コンピューティングはお客様の責任となります。

  • Python UDFs (@udf、spark.udf.register) では、リモートワーカーのバージョンと一致するローカル Python マイナーバージョンが必要です。そうしないと、 で失敗しますPYTHON_VERSION_MISMATCH。組み込み SQL 関数と DataFrame オペレーションでは、Python バージョンの一致は必要ありません。