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
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.
Faça login no nó mestre do cluster EMR on ECS. Para mais informações, consulte Fazer login em um cluster.
-
Execute o comando a seguir para alterar o diretório:
cd /var/log/emr/taihao_exporter -
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 -
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
-
Acesse a página Network Connection.
No painel de navegação à esquerda do console EMR, escolha .
Na página Spark, clique em nome do workspace desejado.
Na página EMR Serverless Spark, clique em Normal Network Connection no painel de navegação à esquerda.
Na página Normal Network Connection, clique em Create Network Connection.
-
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
-
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.

-
Adicione uma regra de grupo de segurança.
Na página Clusters, clique em ID do cluster desejado.
Na página Basic Information, clique em link ao lado de Cluster Security Group.
-
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.
ImportanteNã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
Na página do EMR Serverless Spark, clique em Artifacts no painel de navegação à esquerda.
Na página Artifacts, clique em Upload File.
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
Na página do EMR Serverless Spark, clique em Development no painel de navegação à esquerda.
Na aba Development, clique em ícone
.Insira um nome, selecione como tipo de job e clique em OK.
-
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_kafkaNota-
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.
-
Clique em Publish.
Na caixa de diálogo Publish, clique em OK.
-
Inicie o job de streaming.
Clique em Go to O&M.
Clique em START.
Etapa 7: Visualizar os logs
Clique em aba Log Exploration.
-
Na aba Log Exploration, visualize os detalhes e resultados da execução da aplicação.

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.