Todos os produtos
Search
Central de documentação

E-MapReduce:PySpark streaming on EMR Serverless Spark

Última atualização: Jun 27, 2026

O processamento de streams é essencial para análises de big data em tempo real. O EMR Serverless Spark é uma plataforma poderosa e escalável que simplifica o processamento de dados ao eliminar a necessidade de gerenciar servidores. Este tópico demonstra como usar o EMR Serverless Spark para enviar um job de streaming PySpark, destacando sua facilidade de uso e manutenção para processamento de streams.

Pré-requisitos

Crie um workspace. Para mais informações, consulte Criar um workspace.

Procedimento

Etapa 1: Criar um cluster Dataflow e produzir mensagens

  1. Na página EMR on ECS, crie um cluster Dataflow em tempo real que inclua o serviço Kafka. Para mais informações, consulte Criar um cluster.

  2. Faça login no nó mestre do cluster EMR on ECS. Para mais informações, consulte Fazer login em um cluster.

  3. Execute o comando a seguir para alterar o diretório:

    cd /var/log/emr/taihao_exporter
  4. Execute o comando a seguir para criar um tópico:

    # Create a topic named taihaometrics with 10 partitions and a replication factor of 2.
    kafka-topics.sh --partitions 10 --replication-factor 2 --bootstrap-server core-1-1:9092 --topic taihaometrics --create
  5. Execute o comando a seguir para enviar mensagens:

    # Use kafka-console-producer to send messages to the taihaometrics topic.
    tail -f metrics.log | kafka-console-producer.sh --broker-list core-1-1:9092 --topic taihaometrics

Etapa 2: Criar uma conexão de rede

  1. Acesse a página Network Connection.

    1. No painel de navegação à esquerda do console EMR, escolha EMR Serverless > Spark.

    2. Na página Spark, clique em nome do workspace desejado.

    3. Na página EMR Serverless Spark, clique em Normal Network Connection no painel de navegação à esquerda.

  2. Na página Normal Network Connection, clique em Create Network Connection.

  3. Na caixa de diálogo Create Network Connection, configure os parâmetros a seguir e clique em OK.

    Parâmetro

    Descrição

    Name

    Insira um nome para a conexão. Por exemplo, connection_to_emr_kafka.

    VPC

    Selecione a VPC onde seu cluster EMR on ECS está implantado.

    Se nenhuma VPC estiver disponível, clique em Create VPC para acessar o console VPC e criar uma VPC. Para mais informações, consulte Criar e gerenciar uma VPC.

    vSwitch

    Selecione o vSwitch que está na mesma VPC do seu cluster EMR on ECS.

    Caso não haja nenhum vSwitch disponível na zona atual, clique em vSwitch para acessar o console VPC e criar um vSwitch. Para mais informações, consulte Criar e gerenciar um vSwitch.

    Quando o Status for Succeeded, a conexão de rede terá sido criada com sucesso.

Etapa 3: Adicionar uma regra de grupo de segurança

  1. Obtenha o bloco CIDR do vSwitch dos nós do cluster.

    Na página Nodes, clique em nome de um grupo de nós para localizar o vSwitch associado. Em seguida, faça login no console VPC e encontre o bloco CIDR do vSwitch na página vSwitch.

    image

  2. Adicione uma regra de grupo de segurança.

    1. Na página Clusters, clique em ID do cluster desejado.

    2. Na página Basic Information, clique em link ao lado de Cluster Security Group.

    3. Na página Security Group Details, na seção Rules, clique em Add Rule. Configure os parâmetros a seguir e clique em OK.

      Parâmetro

      Descrição

      Source

      Insira o bloco CIDR do vSwitch obtido na etapa anterior.

      Importante

      Não defina este valor como 0.0.0.0/0, pois isso expõe o cluster a acessos externos.

      Destination (current instance)

      Insira a porta 9092.

Etapa 4: Fazer upload de pacotes JAR para o OSS

Extraia o arquivo kafka.zip e envie todos os pacotes JAR contidos nele para o OSS. Para mais informações, consulte Upload simples.

Etapa 5: Fazer upload do arquivo de recurso

  1. Na página do EMR Serverless Spark, clique em Artifacts no painel de navegação à esquerda.

  2. Na página Artifacts, clique em Upload File.

  3. Na caixa de diálogo Upload File, clique em área de upload e selecione o arquivo pyspark_ss_demo.py.

Etapa 6: Criar e iniciar um job de streaming

  1. Na página do EMR Serverless Spark, clique em Development no painel de navegação à esquerda.

  2. Na aba Development, clique em ícone image.

  3. Insira um nome, selecione Application (Streaming) > PySpark como tipo de job e clique em OK.

  4. Na nova aba de desenvolvimento, configure os parâmetros a seguir, mantenha os valores padrão para os demais e clique em Save.

    Parâmetro

    Descrição

    Main Python Resources

    Selecione o arquivo pyspark_ss_demo.py enviado na página Resource Upload na etapa anterior.

    Engine Version

    Selecione a versão do Spark. Para mais informações, consulte Versões do mecanismo.

    Execution Parameters

    Insira o endereço IP interno do nó core-1-1 do cluster. Você encontra esse IP na página Nodes, dentro do grupo de nós Core.

    Spark Configuration

    Especifique as configurações do Spark. Veja um exemplo abaixo.

    spark.jars oss://path/to/commons-pool2-2.11.1.jar,oss://path/to/kafka-clients-2.8.1.jar,oss://path/to/spark-sql-kafka-0-10_2.12-3.3.1.jar,oss://path/to/spark-token-provider-kafka-0-10_2.12-3.3.1.jar
    spark.emr.serverless.network.service.name connection_to_emr_kafka
    Nota
    • spark.jars: Os caminhos no OSS dos pacotes JAR externos necessários. Substitua os caminhos de exemplo pelos caminhos no OSS dos JARs enviados na Etapa 4.

    • spark.emr.serverless.network.service.name: O nome da conexão de rede. Substitua o valor de exemplo pelo nome da sua conexão de rede criada na Etapa 2.

  5. Clique em Publish.

  6. Na caixa de diálogo Publish, clique em OK.

  7. Inicie o job de streaming.

    1. Clique em Go to O&M.

    2. Clique em START.

Etapa 7: Visualizar os logs

  1. Clique em aba Log Exploration.

  2. Na aba Log Exploration, visualize os detalhes e resultados da execução da aplicação.

    image

Tópicos relacionados

Para ver um exemplo de fluxo de trabalho de desenvolvimento com PySpark, consulte Guia de início rápido de desenvolvimento com PySpark.