View a markdown version of this page

Execute sessões interativas com o Amazon EMR no EKS por meio do Spark Connect - Amazon EMR

As traduções são geradas por tradução automática. Em caso de conflito entre o conteúdo da tradução e da versão original em inglês, a versão em inglês prevalecerá.

Execute sessões interativas com o Amazon EMR no EKS por meio do Spark Connect

Com o Amazon EMR na versão EKS emr-7.14.0 e posterior (ou emr-spark-8.1.0 posterior), você pode se conectar a um endpoint gerenciado do Spark Connect a partir de PySpark clientes autogerenciadosVS Code, como, PyCharm e notebooks. Jupyter O Spark Connect usa uma arquitetura cliente-servidor que separa o código do aplicativo do processo do driver do Spark. Você desenvolve e depura PySpark código em seu IDE local enquanto as operações do Spark são executadas em seu cluster EKS por meio do Amazon EMR no EKS. O Spark Connect oferece os seguintes benefícios:

  • Conecte-se ao Amazon EMR no EKS a partir de qualquer PySpark cliente, incluindo VS CodePyCharm, e notebooks Jupyter

  • Defina pontos de interrupção e percorra o PySpark código em seu IDE enquanto DataFrames executa remotamente dados em escala de produção

  • Execute em seu próprio cluster EKS com controle total sobre a configuração de computação, rede e segurança

Um endpoint Spark Connect é um endpoint gerenciado em seu cluster virtual que hospeda um servidor Spark Connect. Quando você cria um endpoint Spark Connect, o Amazon EMR no EKS provisiona um driver Spark com um servidor gRPC em seu cluster EKS. Seu PySpark cliente local envia operações DataFrame de SQL para o driver por meio do endpoint gRPC. Para interagir com o endpoint, você obtém as credenciais da sessão usando a GetManagedEndpointSessionCredentials API. Cada endpoint oferece suporte a várias sessões simultâneas.

Pré-requisitos

Antes de criar um endpoint do Spark Connect, verifique se você tem o seguinte.

Os requisitos a seguir são específicos dos endpoints do Spark Connect:

  • Um cluster EKS com pelo menos uma sub-rede privada. Para rotear o tráfego gRPC para o endpoint, o Spark Connect provisiona um Network Load Balancer (NLB) interno, que exige uma sub-rede privada em sua VPC. O NLB é interno e não está exposto à Internet pública.

  • O controlador do balanceador de AWS carga instalado em seu cluster EKS, para que o Amazon EMR no EKS possa provisionar o NLB interno. Para obter instruções, consulte Instalando o controlador do balanceador de AWS carga.

Você também precisa dos seguintes recursos padrão do Amazon EMR em EKS. Se você já executa trabalhos no Amazon EMR no EKS, provavelmente tem os seguintes:

Permissões obrigatórias

Adicione as seguintes permissões à sua função do IAM para criar e interagir com um endpoint do 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" ] } } } ] }

Criar uma configuração de segurança

Os endpoints do Spark Connect exigem uma configuração de segurança associada ao seu cluster virtual. A configuração de segurança define as configurações de autenticação e autorização para o endpoint e especifica o namespace do sistema em que a infraestrutura do Spark Connect é executada.

nota

Cada configuração de segurança tem um relacionamento individual com um cluster virtual. Você não pode reutilizar a mesma configuração de segurança em vários clusters virtuais.

Crie uma configuração de segurança com um namespace do 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" } } }'

namespaceNa configuração de segurança está o namespace do sistema em que os componentes da infraestrutura do Spark Connect são executados. Isso é separado do namespace do usuário especificado ao criar o cluster virtual.

Crie um cluster virtual com o Spark Connect ativado

Crie um cluster virtual associado à configuração de segurança. Você deve sessionEnabled definir true e fornecer o securityConfigurationId da etapa anterior.

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 esse cluster virtual à configuração de segurança criada na etapa anterior. Isso é necessário para os endpoints do Spark Connect.

--session-enabled— Permite o suporte ao endpoint Spark Connect no cluster virtual. Sem esse sinalizador, você não pode criar endpoints do Spark Connect nesse cluster virtual.

Crie um endpoint Spark Connect

Depois de criar uma configuração de segurança, crie um endpoint gerenciado pelo Spark Connect no seu cluster virtual.

nota

O primeiro endpoint Spark Connect em um cluster EKS leva mais tempo para se tornar ACTIVE do que os subsequentes. Para o primeiro endpoint no cluster EKS, o Amazon EMR on EKS provisiona os componentes de rede compartilhados usados para conectividade — um Network Load Balancer (NLB) interno, um endpoint de interface VPC () e uma implantação única do roteador Envoy (no spark-connect-router namespace AWS PrivateLink) que roteia o tráfego gRPC para endpoints — o que pode levar vários minutos. Esses componentes são criados somente uma vez por cluster EKS e são reutilizados por todos os endpoints posteriores nesse cluster, para que os endpoints subsequentes sejam inicializados mais rapidamente.

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

Por padrão, um endpoint gerenciado pelo Spark Connect começa com 2 executores. Se você não fornecer substituições de configuração, o endpoint usará esse padrão. Para executar com um número diferente de executores, spark.executor.instances defina as substituições de configuração, conforme mostrado no exemplo a seguir.

O exemplo a seguir inclui substituições de configuração do Spark com alocação dinâmica e uma configuração de monitoramento para exportar registros do Spark para o 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" } }'

Monitore o status do endpoint:

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

Espere até que o estado do endpoint esteja ACTIVE antes de se conectar.

Conecte-se a um endpoint Spark Connect

Depois que o endpoint estiver ativo, obtenha as credenciais da sessão e conecte-se a partir de um PySpark cliente.

Para se conectar a um endpoint do Spark Connect
  1. Obtenha o URL do proxy de autenticação na descrição do endpoint:

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

    A resposta inclui o authProxyUrl campo quando o endpoint éACTIVE.

  2. Obtenha um token de sessão para o 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"

    A resposta inclui um token de sessão:

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

    Use o credentials.token valor como x-aws-proxy-auth parâmetro ao se conectar ao URL do proxy de autenticação.

  3. Instale o PySpark cliente correspondente à versão do Spark em seu endpoint (Spark 3.5.8 paraemr-7.14.0, Spark 4.1.1 para): 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. Conecte-se PySpark usando o URL do proxy de autenticação e o token de sessão:

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

O script Python a seguir obtém a URL do proxy de autenticação e o token de sessão e, em seguida, se conecta ao endpoint do 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()

Considerações e limitações

Considere o seguinte ao executar cargas de trabalho interativas por meio do Spark Connect no Amazon EMR no EKS.

  • O Spark Connect é compatível com o Amazon EMR na versão EKS emr-7.14.0 e posterior ou emr-spark-8.1.0 posterior.

  • O Spark Connect suporta DataFrame e inclui APIs SQL. PySpark O Spark Connect não é compatível com RDD-based APIs.

  • Os tokens de sessão são limitados no tempo. Quando um token expira, as chamadas gRPC falham com um erro de autenticação. Ligue GetManagedEndpointSessionCredentials para obter um novo token e crie um novo SparkSession com o token atualizado.

  • Cada configuração de segurança tem um relacionamento individual com um cluster virtual. Você não pode compartilhar uma configuração de segurança em vários clusters virtuais.

  • Você deve excluir todos os endpoints usando uma configuração de segurança antes de excluir a configuração de segurança.

  • A PySpark versão instalada localmente deve corresponder à versão do Apache Spark em seu endpoint (Spark 3.5.8 para, Spark 4.1.1 paraemr-7.14.0). emr-spark-8.1.0 Uma incompatibilidade de versão causa erros de conexão ou comportamento inesperado.

  • O tipo de endpoint do Spark Connect é. SPARK_CONNECT Isso é diferente dos terminais interativos Livy (tipoJUPYTER_ENTERPRISE_GATEWAY).

  • O sessionIdleTimeoutInMinutes parâmetro controla por quanto tempo uma sessão ociosa persiste antes do encerramento automático. O padrão é de 60 minutos.

  • Os endpoints do Spark Connect não oferecem suporte à propagação de identidade confiável.

  • Os endpoints do Spark Connect ainda não oferecem suporte ao controle de acesso refinado (FGAC) do Lake Formation. Para aplicar o controle de acesso, use a função de execução do IAM associada ao endpoint.

  • Os endpoints do Spark Connect usam um Network Load Balancer (NLB) para rotear o tráfego gRPC. O NLB é criado quando o primeiro endpoint do Spark Connect é criado e só é excluído quando o último cluster virtual habilitado para a sessão é excluído. O Amazon EMR no EKS também executa um único roteador Envoy (no spark-connect-router namespace, um por cluster EKS) que roteia o tráfego gRPC para endpoints; ele é criado com o primeiro endpoint Spark Connect e encerrado quando o último cluster virtual habilitado para sessão é excluído. Você é responsável pelos custos de NLB e pela computação EKS consumida pelo roteador Envoy enquanto eles existem, além dos recursos de computação EKS consumidos pelo driver e pelos executores do Spark durante sua sessão.

  • Os UDFs do Python (@udf,spark.udf.register) exigem que a versão secundária local do Python corresponda à versão do trabalhador remoto, ou falham nela. PYTHON_VERSION_MISMATCH Built-in As funções e DataFrame operações SQL não exigem uma correspondência de versão do Python.