Chapitre 4 Spark
Big Data - 2019 1
Plan
Critique de MapReduce
Apache Spark
Fonctionnalités de Spark
Ecosystème de Spark
Architecture de Spark
RDD
DataFrame
Dataset
2
Critique de MapReduce
Si on souhaite mettre en place une solution complexe, il est nécessaire d’enchaîner une série de jobs MapReduce et de les exécuter séquentiellement.
=> Il est difficile d’exprimer des opérations complexes en n’utilisant que MapReduce.
3
Critique de MapReduce
Après une opération Map ou Reduce, le résultat est écrit sur disque.
Les Mappers et Reducers communiquent entre eux à travers les
données écrites sur disque.
=> Opérations de lectures/écritures très coûteuses en temps !
=> MapReduce est une bonne solution pour les traitements à
passe unique .
=> MapReduce n’est pas la meilleure solution pour les traitements
à plusieurs passes .
Solution : Spark
4
Spark : Présentation
Framework de conception et d'exécution Map/Reduce.
Originellement (2014) un projet de l'université de Berkeley en Californie, désormais un logiciel libre de la fondationApache.
LicenceApache.
Trois modes d'exécution:
Cluster Spark natif.
Hadoop (YARN).
- C’est un moteur MapReduce plus évolué, plus rapide pour les tâches impliquant de multiples maps et/ou reduce.
Utilisation de la mémoire pour optimiser les traitements.
Des API’s pour faciliter et optimiser les étapesd’analyses.
2
Spark : Plus rapide
Big Data - 2019 3
Apache Spark
Framework de traitement Big Data open source .
Développé en 2009 par AMPLab - University of California
Berkeley.
- En 2010 il est passé open source sous forme de projet
Apache.
7
Apache Spark
Framework complet et unifié adapté au traitement batch et temps réel de divers types de données.
Il permet à des applications sur clusters Hadoop d’être exécutées jusqu’à 100 fois plus vite en mémoire, 10 fois plus vite sur disque.
Il est écrit en Scala et s’exécute sur la Machine Virtuelle Java (JVM).
Il permet d’écrire facilement et rapidement des applications en Java, Scala, R ou Python.
Advertisement
Il est possible de l’utiliser de façon interactive pour requêter les données depuis un shell.
8
Apache Spark
En plus des opérations Map et Reduce, Spark supporte les requêtes SQL et le streaming de données.
Spark propose des fonctionnalités de machine learning et de traitements orientés graphe.
Ces possibilités peuvent être utilisées en stand-alone ou en les combinant en une chaîne de traitement complexe.
9
Apache Spark
Spark permet de développer des pipelines de traitement de données complexes, à plusieurs étapes, en s’appuyant sur des graphes orientés acycliques (DAG).
Il permet de partager les données en mémoire entre les graphes, de façon à ce que plusieurs jobs puissent travailler sur le même jeu de données.
10
Apache Spark
Il est possible de déployer des applications Spark sur un cluster Hadoop v1 existant (avec SIMR - Spark-Inside-MapReduce), ou sur un cluster Hadoop v2 YARN.
En général, Spark n’est pas considéré comme un remplaçant de Hadoop mais comme une alternative au MapReduce d’Hadoop.
11
Apache Spark
Spark maintient les résultats intermédiaires en mémoire plutôt que sur disque : ce qui est très utile en particulier lorsqu’il est nécessaire de travailler plusieurs fois sur le même jeu de données.
Le moteur d’exécution de Spark est conçu pour travailler aussi bien en mémoire que sur disque.
Spark essaie de stocker le plus possible en mémoire avant de basculer sur disque.
Il est capable de travailler avec une partie des données en mémoire, une autre sur disque.
12
Fonctionnalités de Spark
Les autres fonctionnalités proposées par Spark comprennent :
Des fonctions autres que Map et Reduce ;
L’évaluation paresseuse des requêtes, ce qui aide à optimiser le
workflow global de traitement ;
Des APIs concises et cohérentes en Scala, Java et Python ;
Un shell interactif pour Scala et Python (non disponible encore en
Java).
13
Ecosystème de Spark
14
Ecosystème de Spark
Spark Core :
Moteur de base de Spark
Gestion de tâches
Gestion de mémoire
Récupération d’erreurs
Interaction avec le stockage
Définition des classes de RDD
15
Ecosystème de Spark
A part les API principales de Spark, son écosystème contient des librairies additionnelles telles que :
Spark streaming : est utilisé pour le traitement temps-réel des données en flux. Il s’appuie sur un mode de traitement en "micro batch" et utilise Dstream (c’est-à-dire une série de RDD (Resilient Distributed Dataset)).
Spark SQL : Grâce à des requêtes de type SQL, Spark SQL permet d’extraire, transformer et charger des données sous différents formats (JSON, Parquet, base de données) et les exposer pour des requêtes ad-hoc.
16
Ecosystème de Spark
MLlib : c’est une librairie de machine learning qui contient tous les algorithmes et utilitaires d’apprentissage classiques, comme la classification, la régression, le clustering, en plus des primitives d’optimisation sous-jacentes.
GraphX : c’est une API (en version alpha) pour les traitements de graphes. Elle étend les RDD de Spark en introduisant le Resilient Distributed Dataset Graph, un multi-graphe orienté avec des propriétés attachées aux nœuds et aux arrêtes. GraphX inclut une collection toujours plus importante d’algorithmes et de builders pour simplifier les tâches d’analyse de graphes.
17
Ecosystème de Spark
Il existe aussi des adaptateurs pour intégration à d’autres produits comme Cassandra (Spark Cassandra Connector) et R (SparkR).
Avec le connecteur Cassandra, il est possible d’utiliser Spark pour accéder à des données stockées dans Cassandra et réaliser des analyses sur ces données.
Advertisement
18
Architecture de Spark
L’architecture de Spark est basée sur trois composants principaux :
19
Architecture de Spark
Spark utilise le système de fichiers HDFS pour le stockage des données. Il peut fonctionner avec n’importe quelle source de données compatible avec Hadoop, dont HDFS, HBase, Cassandra.
L’API permet aux développeurs de créer des applications Spark en utilisant une API standard. L’API existe en Scala, Java et Python.
Spark peut être déployé comme un serveur autonome ou sur un framework de traitements distribués comme Mesos ou YARN.
20
Spark : Facile à utiliser (1/2)
Spark est développé en Scala et supporte quatre langages: Scala, Java, Python (PySpark), R (SparkR).
Une liste d’Operaters pour faciliter la manipulation des données au travers des RDD’S.
Map, filter, groupBy, sort, join, leftOuterJoin, rightOuterJoin, reduce, count, reduceByKey,
groupByKey, first, union, cross, sample, cogroup, take, partionBy, pipe, save,…
Spark : Facile à utiliser (2/2)
Spark est capable de déterminer quand il aura besoin de sérialiser les données / les réorganiser; et ne le fait que quand c'est nécessaire.
On peut également explicitement lui demander de conserver des donnéesen RAM, parce qu'on sait qu'elles seront nécessaires entre plusieurs instances d'écriture disque.
Big Data - 2019 2 4
Spark : Architecture - RDD
Au centre du paradigme employé par Spark, on trouve la notion de RDD, pour Resilient Distributed Datasets.
Il s'agit de larges hashmaps stockées en mémoire et sur lesquelles on peut appliquer des traitements.
Ils sont:
Distribués .
Partitionnés (pour permettre à plusieurs noeuds de traiter les données).
Redondés (limite le risque de perte de données).
En lecture seule : un traitement appliqué à un RDD donne lieu à la création d'un nouveau RDD.
- Deux types d'opérations possibles sur les RDDs:
Une transformation : une opération qui modifie les données d'un RDD. Elle donne lieu à la
création d'un nouveau RDD. Les transformations fonctionnent en mode d'évaluation lazy : elles ne sont exécutées que quand on a véritablement besoin d'accéder aux données. "map" est un exemple de transformation.
Une action : elles accèdent aux données d'un RDD, et nécessitent donc son évaluation
(toutes les transformations ayant donné lieu à la création de ce RDD sont exécutées l'une aprés l'autre). "saveAsTextFile" (qui permet de sauver le contenu d'un RDD) ou "count" (qui renvoie le nombre d'éléments dans un RDD) sont des exemples d'actions.
Spark : Architecture - Exécution
Les applications Spark s’exécutent comme un ensemble de processus indépendants sur un cluster, coordonnés par un objet SparkContext du programme principal, appelé Driver Program.
Pour s’exécuter sur un cluster, le SparkContext se connecte à un Cluster Manager, qui peut être soit un gestionnaire standalone de Spark, soit YARN ou Mesos, pour l’allocation de ressources aux applictions.
Une fois connecté, Spark lance des executors sur les nœuds du cluster, des processus qui lancent des traitements et stockent les données pour les applications.
Il envoie ensuite le code de l’application (Jar ou ficher python) aux executors . SparkContext envoie ensuite les Tasks aux executors pour qu’ils les lancent.
Architecture de Spark
Driver : effectue la distribution les données et le traitement sur les noeuds de travail.
SparkContext : objet du programme principal qui coordonne entre les différents processus sur un cluster.
Executor : processus qui lance des traitements et stocke les données pour les applications
27
Resilient Distributed Dataset (RDD)
Un concept au cœur du framework Spark.
Collection d’éléments découpables par Spark et
traitables de manière distribuée.
- Il peut porter tout type de donnée et est stocké par
Spark sur différentes partitions.
28
Resilient Distributed Dataset
Les RDD supportent deux types d’opérations :
Les transformations :
Advertisement
Ne retournent pas de valeur
Elles retournent un nouveau RDD.
Rien n’est évalué lorsque l’on fait appel à une fonction de
transformation, cette fonction prend juste un RDD et retourne un nouveau RDD.
- Exemples de fonctions :
map,filter,flatMap,
groupByKey,reduceByKey, aggregateByKey…
29
Resilient Distributed Dataset
Les actions :
Evaluent et retournent une nouvelle valeur.
Au moment où une fonction d’action est appelée sur un
objet RDD, toutes les requêtes de traitement des données sont calculées et le résultat est retourné.
- Exemples d’actions :
reduce,collect,count,
first, take, countByKey et foreach.
30
Resilient Distributed Dataset
- Les RDD permettent de réarranger les calculs et d’optimiser
le traitement.
- Les RDD sont reconstructibles : ils sont tolérants aux
pannes car un RDD sait comment recréer et recalculer son ensemble de données.
- Les RDD sont immuables . Pour obtenir une modification
d’un RDD, il faut y appliquer une transformation, qui retournera un nouveau RDD, l’original restera inchangé.
31
Traitement dans Spark
- Spark supporte les « évaluations paresseuses » des requêtes
c’est-à-dire que les transformations ne s’exécutent sur le cluster que si on en a besoin (une action est invoquée).
- Il est possible de demander la persistance d’un RDD :
chargement en mémoire du RDD pour le réutiliser en cas de besoin au lieu de refaire la transformation.
32
DataFrames
2011 RDD :
- Collection distribuée
- Opérateurs fonctionnels
- Ne suis aucun schéma
2013 nouvelle abstraction : DataFrame
- Plus structurée
- Représentation interne plus optimisée que le RDD
- Abstraction principale de SparkSQL
33
Dataset
2015 Spark 2 : Dataset
- Nouvelle abstraction de données plus large que le
DataFrame
- Typé : Dataset [ T ]
- Avantage : travailler avec des expressions typées dont on
connait toutes les propriétés
34
Spark : Un framework analytique
BIG DATA - 2019
35
Spark : Avantages
Performances supérieures à celles de Hadoop pour une large quantité de problèmes; et presque universellement au moins équivalentes pour le reste.
API simple et bien documentée; très simple à utiliser. Paradigme plus souple qui permet un développement conceptuellement plus simple.
Très intégrable avec d'autres solutions; peut très facilement lire des données depuis de nombreuses sources, et propose des couches d'interconnexion très faciles à utiliser pour le reste (API dédiée Spark Streaming, Spark SQL).
APIs dédiées pour le traitement de problèmes en machine learning (Spark MLlib) et graphes (Spark GraphX).
BIG DATA - 2019
36
Spark : Performances
Advertisement
Fortement dépendantes du problème mais d'une manière générale supérieures à Hadoop. Dans le cas de problèmes complexes, effectuant de nombreuses opérations sur les données et notamment sur les mêmes données antérieures, fortement supérieures (jusqu'à 100x plus rapide).
- Même sans travailler plusieurs fois surles
mêmes données ou sans persistance explicite particulière, le paradigme de programme parallélisé en graphe acyclique couplé aux RDDs et leurs propriétés permet des améliorations notables (~10x) pour de nombreux problèmes dépassant le cadre rigide du simple map/shuffle/reduce lancé une fois.
BIG DATA - 2019
Spark : Inconvénients
Spark consomme beaucoup plus de mémoire vive que Hadoop, puisqu'il est susceptible de garder une multitude de RDDs en mémoire. Les serveurs nécessitent ainsi plus de RAM.
Il est moins mature que Hadoop.
Son cluster manager (« Spark Master ») est encore assez immature et laisse à désirer en terme de déploiement / haute disponibilité / fonctionnalités additionnelles du même type; dans les faits, il est souvent déployé via Yarn, et souvent sur un cluster Hadoop existant.
BIG DATA - 2019
Spark : Usage /Manipulation (1/5)
BIG DATA - 2019
39
Spark : Usage /Manipulation (2/5)
BIG DATA - 2019
40
Spark : Usage /Manipulation (3/5)
BIG DATA - 2019
41
Spark : Usage /Manipulation (4/5)
BIG DATA - 2019
42
Spark : Usage /Manipulation (5/5)
La manipulation de l’exemple "Word Count" est fournie en "TP Spark" sous la VM Cloudera où Spark est installé et prêt à l’utilisation.
BIG DATA - 2019
43
Exemples 2 - TP
Construction d’un dataframe Spark appelé « restaurants_df » à
partir de la table restaurant existante dans Cassandra. Le
schéma (noms des colonnes) est connu mais le type des
colonnes n’est pas connu au sein du dataframe.
val restaurants_df =
spark.read.cassandraFormat("restaurant",
"resto_ny").load()
restaurants_df.printSchema()
restaurants_df.show()
44
Exemples
Réalisation d’un filtre correspondant au where dans SQL :
val manhattan = restaurants_df.filter("borough =
'MANHATTAN'")
manhattan.show()
45
Exemples
Construction d’un Dataset dont les colonnes sont typées
case class Restaurant(id: Integer, Name: String,
borough: String, BuildingNum: String, Street:
String, ZipCode: Integer, Phone: String,
CuisineType: String)
val restaurants_ds= restaurants_df.as
Pour ceci il a fallu définir une classe dans le langage de programmation (ici, Scala) et demander la conversion.
46
Exemples
Réalisation du même filtre sur le dataset :
val r = restaurants_ds.filter(r => r.borough ==
"MANHATTAN")
Réalisation d’agrégats par arrondissement :
val comptage_par_borough =
restaurants_ds.groupBy("borough").count()
47
Conclusion
RDD : données au schéma très flexible, mais beaucoup plus difficile à manipuler.
DataFrame/ Dataset :
Données au schéma très contraint offrant un niveau de sécurité élevé.
Concepteur : possibilité de référencer des champs et de leur appliquer des opérations standards en fonction de leur type sans avoir à écrire une fonction spécifique pour la moindre opération
=> Rend le code beaucoup lisible et concis.
- Système : la connaissance du schéma facilite les contrôles avant exécution
48