Le shell Spark est un environnement interactif permettant d'explorer les données et de se familiariser avec l'API Spark. Cette rubrique explique comment démarrer le shell Spark sur un cluster E-MapReduce (EMR) et manipuler les resilient distributed datasets (RDD).
Prérequis
Avant de commencer, assurez-vous de disposer des éléments suivants :
Un cluster EMR avec Spark installé
Un accès SSH au nœud maître du cluster. Consultez Connexion à un cluster
Démarrage du shell Spark
Le shell Spark prend en charge Scala et Python.
Connectez-vous au nœud maître via SSH.
-
Démarrez le shell :
spark-shellUne fois le shell démarré, un SparkContext est disponible sous la forme
sc. Tout SparkContext créé manuellement sera ignoré ; utilisez lescprécréé. Utilisez l'option--masterpour spécifier l'URL du maître du cluster et l'option--jarspour ajouter des packages JAR au classpath. Séparez plusieurs chemins JAR par des virgules. Pour obtenir la liste complète des options, exécutezspark-shell --help.
Notions de base sur les RDD
Un RDD est une collection d'éléments que Spark peut traiter en parallèle. Il existe deux méthodes de création des RDD :
-
Collection parallélisée : encapsulez une collection en mémoire avec
sc.parallelize():val data = Array(1, 2, 3, 4, 5) val distData = sc.parallelize(data) -
Jeu de données externe : lisez un fichier depuis un système de fichiers tel que Hadoop Distributed File System (HDFS), HBase, un système de fichiers partagé ou toute source de stockage fournissant un Hadoop InputFormat :
val distFile = sc.textFile("data.txt")
Transformations et actions
Les opérations sur les RDD se divisent en deux catégories :
Transformations : elles définissent un nouveau RDD à partir d'un RDD existant. Les transformations sont paresseuses (lazy) : Spark enregistre les calculs à effectuer mais n'exécute rien tant qu'une action n'est pas appelée.
Actions : elles déclenchent le calcul et renvoient un résultat au pilote ou écrivent les données dans un stockage.
Comprendre cette distinction est essentiel en pratique. Lorsque vous enchaînez les transformations, Spark construit un plan de calcul. Ce n'est qu'à l'appel d'une action que Spark décompose ce plan en tâches et les exécute sur le cluster.
Transformations
| Opération | Description |
|---|---|
map() |
Applique une fonction à chaque élément et renvoie un nouveau RDD. |
flatMap() |
Applique une fonction à chaque élément, aplatit les résultats et renvoie un nouveau RDD. |
filter() |
Renvoie un nouveau RDD contenant uniquement les éléments pour lesquels la fonction renvoie true. |
distinct() |
Renvoie un nouveau RDD sans doublons. |
union() |
Renvoie un nouveau RDD contenant tous les éléments des deux RDD. |
intersection() |
Renvoie un nouveau RDD contenant uniquement les éléments présents dans les deux RDD. |
subtract() |
Renvoie un nouveau RDD dont les éléments du second RDD ont été retirés. |
cartesian() |
Renvoie un nouveau RDD correspondant au produit cartésien des deux RDD. |
Actions
| Opération | Description |
|---|---|
collect() |
Renvoie tous les éléments du RDD au pilote. |
count() |
Renvoie le nombre d'éléments dans le RDD. |
countByValue() |
Renvoie le nombre d'occurrences de chaque élément dans un RDD. |
reduce() |
Agrège tous les éléments d'un RDD, par exemple pour calculer la somme. |
fold(0)(func) |
Similaire à reduce(), mais utilise une valeur zéro comme valeur initiale. |
aggregate(0)(seqOp, combOp) |
Similaire à reduce(), mais le type de retour peut différer du type des éléments. |
foreach(func) |
Applique une fonction à chaque élément d'un RDD. |
Exemple : calcul de la longueur totale des lignes
Cet exemple lit un fichier texte depuis HDFS, utilise une transformation map pour obtenir la longueur de chaque ligne, puis utilise une action reduce pour additionner ces longueurs.
Étape 1 : Téléchargement du fichier de données
Connectez-vous au nœud maître via SSH.
-
Créez un fichier nommé
data.txtavec le contenu suivant :Hello Spark This is a test file 1234567890 -
Téléchargez le fichier vers HDFS :
hadoop fs -put data.txt /user/root/
Étape 2 : Démarrage du shell Spark
spark-shell
Étape 3 : Calcul de la longueur totale des lignes
Exécutez les commandes suivantes dans le shell :
val lines = sc.textFile("data.txt")
val lineLengths = lines.map(s => s.length)
val totalLength = lineLengths.reduce((a, b) => a + b)
Voici le détail de chaque ligne :
val lines = sc.textFile("data.txt"): crée un RDD de base pointant vers le fichier. Le fichier n'est pas encore lu ;linesn'est qu'un pointeur.val lineLengths = lines.map(s => s.length): définit une transformation qui calculera la longueur de chaque ligne. Rien n'est exécuté à ce stade.val totalLength = lineLengths.reduce((a, b) => a + b): il s'agit d'une action. Spark lit le fichier, applique la transformationmapet exécute la réduction sur le cluster.