Tous les produits
Search
Centre de documentation

E-MapReduce:Exemple de traitement Spark par lots

Dernière mise à jour :Aug 09, 2026

Utilisez Spark pour lire les données de Simple Log Service (SLS) en mode lot (traitement hors ligne), c'est-à-dire comme un ensemble de données borné avec une heure de début et de fin définies, plutôt que comme un flux continu. EMR prend en charge deux approches : Spark RDD et Spark SQL.

Prérequis

Avant de commencer, assurez-vous de disposer des éléments suivants :

  • Un cluster EMR avec Spark installé

  • Un projet Simple Log Service et un Logstore contenant les données à lire

  • Un AccessKey ID et un AccessKey secret pour un utilisateur RAM disposant d'un accès en lecture au Logstore

  • (Spark SQL uniquement) Un accès au nœud du cluster EMR où se trouvent les fichiers JAR

Utiliser Spark RDD pour lire les données depuis Simple Log Service

Exemple de code

Le programme Scala suivant lit les données de journal d'un Logstore sur une plage horaire spécifiée et enregistre la sortie dans un chemin du système de fichiers.

// TestBatchLoghub.scala

object TestBatchLoghub {
  def main(args: Array[String]): Unit = {
    if (args.length < 6) {
      System.err.println(
        """Usage: TestBatchLoghub <sls project> <sls logstore> <sls endpoint>
          |  <access key id> <access key secret> <output path> <start time> <end time=now>
        """.stripMargin)
      System.exit(1)
    }

    val loghubProject    = args(0)
    val logStore         = args(1)
    val endpoint         = args(2)
    val accessKeyId      = args(3)     // Read from environment variable at submit time
    val accessKeySecret  = args(4)     // Read from environment variable at submit time
    val outputPath       = args(5)
    val startTime        = args(6).toLong

    val sc = new SparkContext(new SparkConf().setAppName("test batch loghub"))
    var rdd: JavaRDD[String] = null

    if (args.length > 7) {
      // Read log data between startTime and endTime
      rdd = LoghubUtils.createRDD(sc, loghubProject, logStore, accessKeyId, accessKeySecret, endpoint, startTime, args(7).toLong)
    } else {
      // Read log data from startTime to now
      rdd = LoghubUtils.createRDD(sc, loghubProject, logStore, accessKeyId, accessKeySecret, endpoint, startTime)
    }

    rdd.saveAsTextFile(outputPath)
  }
}

La méthode LoghubUtils.createRDD() renvoie un JavaRDD[String] où chaque élément correspond à une entrée de journal. Fournissez sept arguments pour lire jusqu'à l'heure actuelle, ou huit arguments pour spécifier explicitement une heure de fin.

Pour la configuration Maven POM, consultez aliyun-emapreduce-demo.

Paramètres de connexion

Paramètre Description Exemple
<sls project> Nom du projet SLS my-project
<sls logstore> Nom du Logstore au sein du projet my-logstore
<sls endpoint> Endpoint SLS pour votre région cn-hangzhou.log.aliyuncs.com
<access key id> AccessKey ID de votre utilisateur RAM Utilisez $ALIBABA_CLOUD_ACCESS_KEY_ID
<access key secret> AccessKey secret de votre utilisateur RAM Utilisez $ALIBABA_CLOUD_ACCESS_KEY_SECRET
<output path> Chemin de sortie pour le résultat RDD oss://my-bucket/output/
<start time> Début de la plage horaire 1700000000
<end time> (Facultatif) Fin de la plage horaire. Par défaut, l'heure actuelle est utilisée. 1700003600

Compiler et exécuter

Étape 1 : Compilez le code.

mvn clean package -DskipTests

Le fichier JAR compilé est enregistré dans le répertoire target/shaded/.

Étape 2 : Soumettez la tâche.

Important

Utilisez la paire de clés AccessKey d'un utilisateur RAM plutôt que celle de votre compte Alibaba Cloud. Une paire de clés au niveau du compte accorde l'accès à toutes les opérations API. Pour créer un utilisateur RAM, consultez Create a RAM user. Stockez les identifiants sous forme de variables d'environnement ; ne les codez pas en dur dans vos scripts.

Définissez les variables d'environnement avant la soumission :

export ALIBABA_CLOUD_ACCESS_KEY_ID=<access_key_id>
export ALIBABA_CLOUD_ACCESS_KEY_SECRET=<access_key_secret>

Soumettez ensuite la tâche :

spark-submit \
  --master yarn-cluster \
  --executor-cores 2 \
  --executor-memory 1g \
  --driver-memory 1g \
  --num-executors 2 \
  --class x.x.x.TestBatchLoghub xxx.jar \
  <sls project> <sls logstore> <sls endpoint> \
  $ALIBABA_CLOUD_ACCESS_KEY_ID $ALIBABA_CLOUD_ACCESS_KEY_SECRET \
  <output path> <start time> [<end time>]

Remplacez x.x.x.TestBatchLoghub par le nom de classe complet réel et xxx.jar par le chemin d'accès réel au fichier JAR. Ajustez les paramètres --executor-cores, --executor-memory et --num-executors en fonction du volume de données et de la capacité de votre cluster.

Utiliser Spark SQL pour lire les données depuis Simple Log Service

Spark SQL utilise la source de données loghub, incluse dans le JAR d'extension EMR Spark.

Lancer spark-sql avec le JAR LogHub

Important

Utilisez la paire de clés AccessKey d'un utilisateur RAM plutôt que celle de votre compte Alibaba Cloud. Stockez les identifiants sous forme de variables d'environnement transmises via --hiveconf ; ne les codez pas en dur dans vos scripts.

spark-sql \
  --jars /opt/apps/SPARK-EXTENSION/spark-extension-current/spark3-emrsdk/* \
  --hiveconf accessKeyId=$ALIBABA_CLOUD_ACCESS_KEY_ID \
  --hiveconf accessKeySecret=$ALIBABA_CLOUD_ACCESS_KEY_SECRET

Si votre cluster EMR utilise Spark 2, remplacez spark3 par spark2 dans le chemin d'accès au JAR :

/opt/apps/SPARK-EXTENSION/spark-extension-current/spark2-emrsdk/*

Créer une table et interroger les données

CREATE TABLE test_sls
USING loghub
OPTIONS (
  endpoint         = 'cn-hangzhou-intranet.log.aliyuncs.com',
  access.key.id    = '${hiveconf:accessKeyId}',
  access.key.secret= '${hiveconf:accessKeySecret}',
  sls.project      = 'test_project',
  sls.store        = 'test_store',
  startingoffsets  = 'earliest'
);

SELECT * FROM test_sls;

Paramètres de connexion

Paramètre Obligatoire Description Exemple
endpoint Oui Endpoint SLS pour votre région cn-hangzhou-intranet.log.aliyuncs.com
access.key.id Oui AccessKey ID, transmis via --hiveconf ${hiveconf:accessKeyId}
access.key.secret Oui AccessKey secret, transmis via --hiveconf ${hiveconf:accessKeySecret}
sls.project Oui Nom du projet SLS test_project
sls.store Oui Nom du Logstore test_store
startingoffsets Non Position de lecture initiale. Utilisez earliest pour lire toutes les données disponibles depuis le début du Logstore. earliest

Utiliser le JAR LogHub dans un environnement de développement local

Pour développer et tester localement avec Spark 3 (les étapes sont identiques pour Spark 2), installez le JAR de source de données EMR dans votre référentiel Maven local.

Étape 1 : Téléchargez le JAR depuis votre cluster EMR.

Copiez le JAR depuis le chemin suivant sur le nœud de votre cluster vers votre machine locale :

/opt/apps/SPARK-EXTENSION/spark-extension-current/spark3-emrsdk/emr-datasources_shaded_2.12

Étape 2 : Installez le JAR dans votre référentiel Maven local.

mvn install:install-file \
  -DgroupId=com.aliyun.emr \
  -DartifactId=emr-datasources_shaded_2.12 \
  -Dversion=3.0.2 \
  -Dpackaging=jar \
  -Dfile=<path-to-downloaded-jar>

Remplacez <path-to-downloaded-jar> par le chemin local où vous avez enregistré le JAR.

**Étape 3 : Ajoutez la dépendance à votre fichier pom.xml.**

<dependency>
  <groupId>com.aliyun.emr</groupId>
  <artifactId>emr-datasources_shaded_2.12</artifactId>
  <version>3.0.2</version>
</dependency>

Références

Pour plus d'informations sur l'utilisation de Spark pour accéder à Kafka, consultez Structured Streaming + Kafka Integration Guide.