View a markdown version of this page

Esegui sessioni interattive con Amazon EMR su EKS tramite Spark Connect - Amazon EMR

Le traduzioni sono generate tramite traduzione automatica. In caso di conflitto tra il contenuto di una traduzione e la versione originale in Inglese, quest'ultima prevarrà.

Esegui sessioni interattive con Amazon EMR su EKS tramite Spark Connect

Con Amazon EMR nella versione EKS emr-7.14.0 e successive (o emr-spark-8.1.0 successive), puoi connetterti a un endpoint Spark Connect gestito da PySpark client autogestiti come, e notebook. VS Code PyCharm Jupyter Spark Connect utilizza un'architettura client-server che separa il codice dell'applicazione dal processo del driver Spark. Sviluppi ed esegui il debug PySpark del codice nel tuo IDE locale mentre le operazioni Spark vengono eseguite sul tuo cluster EKS tramite Amazon EMR su EKS. Spark Connect offre i seguenti vantaggi:

  • Connettiti ad Amazon EMR su EKS da qualsiasi PySpark client VS CodePyCharm, inclusi notebook Jupyter

  • Imposta i punti di interruzione e approfondisci il PySpark codice nel tuo IDE mentre DataFrames esegui i dati su scala di produzione in remoto

  • Esegui sul tuo cluster EKS con il pieno controllo della configurazione di elaborazione, rete e sicurezza

Un endpoint Spark Connect è un endpoint gestito sul tuo cluster virtuale che ospita un server Spark Connect. Quando crei un endpoint Spark Connect, Amazon EMR su EKS esegue il provisioning di un driver Spark con un server gRPC sul tuo cluster EKS. Il PySpark client locale invia DataFrame operazioni SQL al driver tramite l'endpoint gRPC. Per interagire con l'endpoint, si ottengono le credenziali di sessione utilizzando l'API. GetManagedEndpointSessionCredentials Ogni endpoint supporta più sessioni simultanee.

Prerequisiti

Prima di creare un endpoint Spark Connect, assicurati di avere quanto segue.

I seguenti requisiti sono specifici per gli endpoint Spark Connect:

  • Un cluster EKS con almeno una sottorete privata. Per indirizzare il traffico gRPC verso l'endpoint, Spark Connect fornisce un Network Load Balancer (NLB) interno, che richiede una sottorete privata nel tuo VPC. L'NLB è interno e non è esposto alla rete Internet pubblica.

  • Il AWS Load Balancer Controller è installato sul tuo cluster EKS, in modo che Amazon EMR on EKS possa effettuare il provisioning dell'NLB interno. Per istruzioni, consulta Installazione del Load Balancer Controller AWS .

Sono inoltre necessarie le seguenti risorse Amazon EMR on EKS standard. Se esegui già processi su Amazon EMR su EKS, probabilmente li hai già pronti:

Autorizzazioni richieste

Aggiungi le seguenti autorizzazioni al tuo ruolo IAM per creare e interagire con un endpoint 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" ] } } } ] }

Creazione di una configurazione di sicurezza

Gli endpoint Spark Connect richiedono una configurazione di sicurezza associata al cluster virtuale. La configurazione di sicurezza definisce le impostazioni di autenticazione e autorizzazione per l'endpoint e specifica lo spazio dei nomi di sistema in cui viene eseguita l'infrastruttura Spark Connect.

Nota

Ogni configurazione di sicurezza ha una relazione uno-a-uno con un cluster virtuale. Non è possibile riutilizzare la stessa configurazione di sicurezza su più cluster virtuali.

Crea una configurazione di sicurezza con uno spazio dei nomi di sistema:

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" } } }'

namespaceNella configurazione di sicurezza si trova lo spazio dei nomi di sistema in cui vengono eseguiti i componenti dell'infrastruttura Spark Connect. Questo è separato dallo spazio dei nomi utente specificato durante la creazione del cluster virtuale.

Crea un cluster virtuale con Spark Connect abilitato

Crea un cluster virtuale associato alla configurazione di sicurezza. È necessario sessionEnabled impostare true e fornire quanto securityConfigurationId indicato nel passaggio precedente.

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
Importante

--security-configuration-id— Associa questo cluster virtuale alla configurazione di sicurezza creata nel passaggio precedente. Questo è necessario per gli endpoint Spark Connect.

--session-enabled— Abilita il supporto degli endpoint Spark Connect sul cluster virtuale. Senza questo flag, non è possibile creare endpoint Spark Connect su questo cluster virtuale.

Crea un endpoint Spark Connect

Dopo aver creato una configurazione di sicurezza, crea un endpoint gestito da Spark Connect sul tuo cluster virtuale.

Nota

Il primo endpoint Spark Connect su un cluster EKS impiega più tempo a diventare ACTIVE rispetto a quelli successivi. Per il primo endpoint del cluster EKS, Amazon EMR on EKS fornisce i componenti di rete condivisi utilizzati per la connettività: un Network Load Balancer (NLB) interno, un endpoint di interfaccia VPC (AWS PrivateLink) e una distribuzione unica del router Envoy (nel spark-connect-router namespace) che indirizza il traffico gRPC agli endpoint, operazione che può richiedere diversi minuti. Questi componenti vengono creati una sola volta per cluster EKS e vengono riutilizzati da tutti gli endpoint successivi su quel cluster, in modo che gli endpoint successivi si avviino più velocemente.

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
Nota

Per impostazione predefinita, un endpoint gestito da Spark Connect inizia con 2 esecutori. Se non fornisci le sostituzioni di configurazione, l'endpoint utilizza questa impostazione predefinita. Per l'esecuzione con un numero diverso di esecutori, imposta le sostituzioni di configurazione, come mostrato spark.executor.instances nell'esempio seguente.

L'esempio seguente include le sostituzioni della configurazione di Spark con allocazione dinamica e una configurazione di monitoraggio per esportare i log di Spark in 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" } }'

Monitora lo stato dell'endpoint:

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

Attendi lo stato dell'endpoint ACTIVE prima di connetterti.

Connettiti a un endpoint Spark Connect

Dopo che l'endpoint è attivo, ottieni le credenziali di sessione e connettiti da un client. PySpark

Per connetterti a un endpoint Spark Connect
  1. Ottieni l'URL del proxy di autenticazione dalla descrizione dell'endpoint:

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

    La risposta include il authProxyUrl campo in cui si trova l'endpoint. ACTIVE

  2. Ottieni un token di sessione per l'endpoint:

    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"

    La risposta include un token di sessione:

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

    Usa il credentials.token valore come x-aws-proxy-auth parametro quando ti connetti all'URL del proxy di autenticazione.

  3. Installa il PySpark client corrispondente alla versione di Spark sul tuo endpoint (Spark 3.5.8 per, Spark 4.1.1 peremr-7.14.0): emr-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. Connettiti PySpark utilizzando l'URL del proxy di autenticazione e il token di sessione:

    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()

Il seguente script Python ottiene l'URL del proxy di autenticazione e il token di sessione, quindi si connette all'endpoint 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()

Considerazioni e limitazioni

Considera quanto segue quando esegui carichi di lavoro interattivi tramite Spark Connect su Amazon EMR su EKS.

  • Spark Connect è supportato con Amazon EMR nella versione EKS emr-7.14.0 e successive, oppure successive. emr-spark-8.1.0

  • Spark Connect supporta DataFrame e API SQL in. PySpark Spark Connect non supporta le API. RDD-based

  • I token di sessione sono limitati nel tempo. Quando un token scade, le chiamate gRPC hanno esito negativo con un errore di autenticazione. Chiama GetManagedEndpointSessionCredentials per ottenere un nuovo token e crearne uno nuovo SparkSession con il token aggiornato.

  • Ogni configurazione di sicurezza ha una relazione uno-a-uno con un cluster virtuale. Non è possibile condividere una configurazione di sicurezza tra più cluster virtuali.

  • È necessario eliminare tutti gli endpoint utilizzando una configurazione di sicurezza prima di poter eliminare la configurazione di sicurezza.

  • La PySpark versione installata localmente deve corrispondere alla versione di Apache Spark sull'endpoint (Spark 3.5.8 per, Spark 4.1.1 per). emr-7.14.0 emr-spark-8.1.0 Una mancata corrispondenza della versione causa errori di connessione o comportamenti imprevisti.

  • Il tipo di endpoint Spark Connect è. SPARK_CONNECT Questo è diverso dagli endpoint interattivi Livy (tipo). JUPYTER_ENTERPRISE_GATEWAY

  • Il sessionIdleTimeoutInMinutes parametro controlla per quanto tempo una sessione inattiva persiste prima della chiusura automatica. L’impostazione predefinita è 60 minuti.

  • Gli endpoint Spark Connect non supportano la Trusted Identity Propagation.

  • Gli endpoint Spark Connect non supportano ancora il controllo granulare degli accessi (FGAC) di Lake Formation. Per applicare il controllo degli accessi, utilizza il ruolo di esecuzione IAM associato all'endpoint.

  • Gli endpoint Spark Connect utilizzano un Network Load Balancer (NLB) per indirizzare il traffico gRPC. L'NLB viene creato quando viene creato il primo endpoint Spark Connect e viene eliminato solo quando viene eliminato l'ultimo cluster virtuale abilitato alla sessione. Amazon EMR on EKS esegue anche un singolo router Envoy (nello spazio dei spark-connect-router nomi, uno per cluster EKS) che indirizza il traffico gRPC agli endpoint; viene creato con il primo endpoint Spark Connect e terminato quando viene eliminato l'ultimo cluster virtuale abilitato alla sessione. L'utente è responsabile dei costi NLB e del calcolo EKS consumato dal router Envoy finché esistono, oltre alle risorse di calcolo EKS consumate dal driver e dagli esecutori Spark durante la sessione.

  • Le UDF di Python (@udf,spark.udf.register) richiedono che la versione secondaria locale di Python corrisponda alla versione del lavoratore remoto, altrimenti non funzionano. PYTHON_VERSION_MISMATCH Built-in Le funzioni e le DataFrame operazioni SQL non richiedono una corrispondenza della versione di Python.