Spark est une librairie scala permettant d'effectuer du calcul distribué. Quand vous exécutez du code sur votre machine, seule votre machine est capable de l'exécuter.
Avec Spark, nous pouvons répartir les calculs sur un cluster (plusieurs machines).
Spark SQL : spark propose un ensemble d'outils pour développer des scripts en utilisant du sql.
Il est possible d'utiliser "Spark SQL" qui permet de faire des requêtes SQL en python.
code_sql = spark.sql("SELECT * FROM nom_de_la_table")
display(code_sql)
type(code_sql)
Spark DataFrame vers Pandas DataFrame
Pour transformer une DataFrame Spark en Pandas, vous pouvez utiliser la méthode toPandas().
spark.sql("SELECT * FROM nom_de_la_table").toPandas()
Spark SQL par "Méthodes"
Il est possible de faire la même chose en utilisant des fonctions.
Table() : permet de sélectionner une table.
resultat_table = spark.table('nom_de_la_table')
display(resultat_table)
Limit() : permet de limiter le nombre de retours.
resultat_table = spark.table('nom_de_la_table')\
.limit(10)
display(resultat_table)
Select() : permet de sélectionner les colonnes.
resultat_table = spark.table('nom_de_la_table')\
.select("colonne1","colonne2")
display(resultat_table)
Distinct() : permet d'enlever les doublons.
resultat_table = spark.table('nom_de_la_table')\
.select("colonne1","colonne2").distinct()
display(resultat_table)
OrderBy() : permet de trier en fonction d'une ou plusieurs colonnes.
resultat_table = spark.table('nom_de_la_table')\
.select("colonne1","colonne2")\
.distinct()\
.orderBy('colonne1')
display(resultat_table)
Les agrégations : Tout comme en SQL et pandas, vous pouvez utiliser les agrégations avec la méthode groupBy.
resultat_table = spark.table('nom_de_la_table')\
.groupBy("colonne1")\
.sum('colonne2')
display(resultat_table)
Les fonctions sur les colonnes : Pour certaines méthodes nous avons besoin d'effectuer des opérations sur les colonnes.
Filtrer : permet d'effectuer un filtre sur les valeurs.
resultat_table = spark.table('nom_de_la_table')
resultat_filtre = resultat_table.filter(resultat_table ['colonne1']>10)
display( resultat_filtre )
WithColumn : permet d'ajouter une nouvelle colonne ou d'en remplacer une existante.
Sauvegarder dans une table
resultat_table = spark.table('nom_de_la_table')
resultat_table.write.saveAsTable("table_eng")
Sauvegarder dans une vue temporaire
spark.table('nom_de_la_table').createTempView("temp_table")
Les jointures