Tous les produits
Search
Centre de documentation

E-MapReduce:PySpark streaming on EMR Serverless Spark

Dernière mise à jour :Aug 09, 2026

Le traitement de flux est essentiel pour l'analyse du Big Data en temps réel. EMR Serverless Spark est une plateforme puissante et évolutive qui simplifie le traitement des données en éliminant la gestion des serveurs. Cette rubrique explique comment utiliser EMR Serverless Spark pour soumettre une tâche de streaming PySpark, en mettant en avant sa simplicité d'utilisation et sa facilité de maintenance.

Prérequis

Vous avez créé un espace de travail. Pour plus d'informations, consultez Créer un espace de travail.

Procédure

Étape 1 : Créer un cluster Dataflow et produire des messages

  1. Sur la page EMR on ECS, créez un cluster Dataflow en temps réel incluant le service Kafka. Pour plus d'informations, consultez Créer un cluster.

  2. Connectez-vous au nœud maître du cluster EMR on ECS. Pour plus d'informations, consultez Se connecter à un cluster.

  3. Exécutez la commande suivante pour changer de répertoire :

    cd /var/log/emr/taihao_exporter
  4. Exécutez la commande suivante pour créer un topic :

    # 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. Exécutez la commande suivante pour envoyer des messages :

    # 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

Étape 2 : Créer une connexion réseau

  1. Accédez à la page Network Connection.

    1. Dans le volet de navigation de gauche de la console EMR, choisissez EMR Serverless > Spark.

    2. Sur la page Spark, cliquez sur le nom de votre espace de travail cible.

    3. Sur la page EMR Serverless Spark, cliquez sur Normal Network Connection dans le volet de navigation de gauche.

  2. Sur la page Normal Network Connection, cliquez sur Create Network Connection.

  3. Dans la boîte de dialogue Create Network Connection, configurez les paramètres suivants et cliquez sur OK.

    Parameter

    Description

    Name

    Saisissez un nom pour la connexion. Par exemple, connection_to_emr_kafka.

    VPC

    Sélectionnez le VPC où votre cluster EMR on ECS est déployé.

    Si aucun VPC n'est disponible, cliquez sur Create VPC pour accéder à la console VPC et créer un VPC. Pour plus d'informations, consultez Créer et gérer un VPC.

    vSwitch

    Sélectionnez le vSwitch situé dans le même VPC que votre cluster EMR on ECS.

    Si aucun vSwitch n'est disponible dans la zone actuelle, cliquez sur vSwitch pour accéder à la console VPC et créer un vSwitch. Pour plus d'informations, consultez Créer et gérer un vSwitch.

    Lorsque le statut Status indique Succeeded, la connexion réseau est créée.

Étape 3 : Ajouter une règle de groupe de sécurité

  1. Obtenez le bloc CIDR du vSwitch pour les nœuds du cluster.

    Sur la page Nodes, cliquez sur le nom d'un groupe de nœuds pour trouver le vSwitch associé. Connectez-vous ensuite à la console VPC et identifiez le bloc CIDR du vSwitch sur la page vSwitch.

  2. Ajoutez une règle de groupe de sécurité.

    1. Sur la page Clusters, cliquez sur l'ID du cluster cible.

    2. Sur la page Basic Information, cliquez sur le lien situé à côté de Cluster Security Group.

    3. Sur la page Security Group Details, dans la section Rules, cliquez sur Add Rule. Configurez les paramètres suivants et cliquez sur OK.

      Parameter

      Description

      Source

      Saisissez le bloc CIDR du vSwitch obtenu à l'étape précédente.

      Important

      Ne définissez pas cette valeur sur 0.0.0.0/0, car cela exposerait le cluster à un accès externe.

      Destination (current instance)

      Saisissez le port 9092.

Étape 4 : Télécharger les packages JAR vers OSS

Extrayez le fichier kafka.zip et téléchargez tous les packages JAR contenus dans l'archive vers OSS. Pour plus d'informations, consultez Téléchargement simple.

Étape 5 : Télécharger le fichier de ressources

  1. Sur la page EMR Serverless Spark, cliquez sur Artifacts dans le volet de navigation de gauche.

  2. Sur la page Artifacts, cliquez sur Upload File.

  3. Dans la boîte de dialogue Upload File, cliquez sur la zone de téléchargement et sélectionnez le fichier pyspark_ss_demo.py.

Étape 6 : Créer et démarrer une tâche de streaming

  1. Sur la page EMR Serverless Spark, cliquez sur Development dans le volet de navigation de gauche.

  2. Sur l'onglet Development, cliquez sur l'icône image.

  3. Saisissez un nom, choisissez Application (Streaming) > PySpark comme type de tâche, puis cliquez sur OK.

  4. Dans le nouvel onglet de développement, configurez les paramètres suivants, conservez les valeurs par défaut pour les autres, puis cliquez sur Save.

    Parameter

    Description

    Main Python Resources

    Sélectionnez le fichier pyspark_ss_demo.py que vous avez téléchargé sur la page Resource Upload à l'étape précédente.

    Engine Version

    Sélectionnez la version Spark. Pour plus d'informations, consultez Versions du moteur.

    Execution Parameters

    Saisissez l'adresse IP interne du nœud core-1-1 du cluster. Vous pouvez trouver cette adresse IP sur la page Nodes, au sein du groupe de nœuds Core.

    Spark Configuration

    Spécifiez les configurations Spark. Voici un exemple.

    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
    Remarque
    • spark.jars : Les chemins OSS des packages JAR externes requis. Remplacez les chemins d'exemple par les chemins OSS des JAR que vous avez téléchargés à l'étape 4.

    • spark.emr.serverless.network.service.name : Le nom de la connexion réseau. Remplacez la valeur d'exemple par le nom de votre connexion réseau défini à l'étape 2.

  5. Cliquez sur Publish.

  6. Dans la boîte de dialogue Publish, cliquez sur OK.

  7. Démarrez la tâche de streaming.

    1. Cliquez sur Go to O&M.

    2. Cliquez sur START.

Étape 7 : Consulter les journaux

  1. Cliquez sur l'onglet Log Exploration.

  2. Sur l'onglet Log Exploration, consultez les détails et les résultats de l'exécution de l'application.

    -------------------------------------------
    Batch: 0
    -------------------------------------------
    +--------+
    |count(1)|
    +--------+
    |    1938|
    +--------+
    
    -------------------------------------------
    Batch: 1
    -------------------------------------------
    +--------+
    |count(1)|
    +--------+
    |    1948|
    +--------+
    
    -------------------------------------------
    Batch: 2
    -------------------------------------------
    +--------+
    |count(1)|
    +--------+
    |    1958|
    +--------+
    
    -------------------------------------------
    Batch: 3
    -------------------------------------------
    +--------+
    |count(1)|
    +--------+
    |    1971|
    +--------+
    
    -------------------------------------------
    Batch: 4
    -------------------------------------------
    +--------+
    |count(1)|
    +--------+
    |    2022|
    +--------+
    
    -------------------------------------------
    Batch: 5
    -------------------------------------------
    +--------+

Rubriques connexes

Pour un exemple de flux de travail de développement PySpark, consultez Démarrage rapide du développement PySpark.