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
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.
Connectez-vous au nœud maître du cluster EMR on ECS. Pour plus d'informations, consultez Se connecter à un cluster.
-
Exécutez la commande suivante pour changer de répertoire :
cd /var/log/emr/taihao_exporter -
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 -
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
-
Accédez à la page Network Connection.
Dans le volet de navigation de gauche de la console EMR, choisissez .
Sur la page Spark, cliquez sur le nom de votre espace de travail cible.
Sur la page EMR Serverless Spark, cliquez sur Normal Network Connection dans le volet de navigation de gauche.
Sur la page Normal Network Connection, cliquez sur Create Network Connection.
-
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é
-
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.
-
Ajoutez une règle de groupe de sécurité.
Sur la page Clusters, cliquez sur l'ID du cluster cible.
Sur la page Basic Information, cliquez sur le lien situé à côté de Cluster Security Group.
-
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.
ImportantNe 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
Sur la page EMR Serverless Spark, cliquez sur Artifacts dans le volet de navigation de gauche.
Sur la page Artifacts, cliquez sur Upload File.
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
Sur la page EMR Serverless Spark, cliquez sur Development dans le volet de navigation de gauche.
Sur l'onglet Development, cliquez sur l'icône
.Saisissez un nom, choisissez comme type de tâche, puis cliquez sur OK.
-
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_kafkaRemarque-
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.
-
Cliquez sur Publish.
Dans la boîte de dialogue Publish, cliquez sur OK.
-
Démarrez la tâche de streaming.
Cliquez sur Go to O&M.
Cliquez sur START.
Étape 7 : Consulter les journaux
Cliquez sur l'onglet Log Exploration.
-
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.