Didacticiel : Charger et transformer des données à l'aide des DataFrames Apache Spark
Ce tutoriel vous montre comment charger et transformer des données à l'aide de l'API DataFrame Apache Spark Python (PySpark), de l'API DataFrame Apache Spark Scala et de l'API SparkDataFrame SparkR dans Databricks.
Si vous utilisez Databricks Free Edition, sélectionnez l'onglet Python pour tous les exemples de code de ce tutoriel. Free Edition ne prend pas en charge R ou Scala. De plus, Free Edition restreint l’accès Internet sortant, vous devez donc upload le fichier CSV via l’interface utilisateur du Workspace au lieu de le download avec du code. Voir l'Étape 1 pour des instructions détaillées.
À la fin de ce tutoriel, vous comprendrez ce qu’est un DataFrame et serez familiarisé avec les tâches suivantes :
- Python
- Scala
- R
- Définissez des variables et copiez des données publiques dans un volume Unity Catalog
- Créer un DataFrame avec Python
- Charger des données dans un DataFrame à partir d'un fichier CSV.
- Afficher et interagir avec un DataFrame
- Enregistrer le DataFrame
- Exécuter des queries SQL dans PySpark
Voir aussi référence de l’API Apache Spark PySpark.
- Définissez des variables et copiez des données publiques dans un volume Unity Catalog
- Créer un DataFrame avec Scala
- Charger des données dans un DataFrame à partir d'un fichier CSV.
- Afficher et interagir avec un DataFrame
- Enregistrer le DataFrame
- Exécuter des requêtes SQL dans Apache Spark.
Voir aussi référence de l'API Scala Apache Spark.
- Définissez des variables et copiez des données publiques dans un volume Unity Catalog
- Créer un SparkDataFrames SparkR
- Charger des données dans un DataFrame à partir d'un fichier CSV.
- Afficher et interagir avec un DataFrame
- Enregistrer le DataFrame
- Exécuter des requêtes SQL dans SparkR
Voir aussi Référence de l'API Apache SparkR.
Qu'est-ce qu'un DataFrame ?
Un DataFrame est une structure de données étiquetée bidimensionnelle avec des colonnes de types potentiellement différents. Vous pouvez considérer un DataFrame comme une feuille de calcul, une table SQL ou un dictionnaire d'objets de série. Les DataFrames Apache Spark offrent un riche ensemble de fonctions (sélectionner des colonnes, filtrer, joindre, agréger) qui vous permettent de résoudre efficacement les problèmes courants d'analyse de données.
Les DataFrames Apache Spark sont une abstraction construite sur les Resilient Distributed Datasets (RDD). Spark DataFrames et Spark SQL utilisent un moteur de planification et d'optimisation unifié, ce qui vous permet d'obtenir des performances quasiment identiques dans toutes les langues prises en charge sur Databricks (Python, SQL, Scala et R).
Exigences
Pour terminer le tutoriel suivant, vous devez respecter les exigences suivantes :
-
Pour utiliser les exemples de ce tutoriel, votre workspace doit avoir Unity Catalog activé. Les workspaces Databricks Free Edition et d'essai gratuit ont Unity Catalog activé par default.
-
Les exemples de ce tutoriel utilisent un volume Unity Catalog pour stocker des exemples de données. Pour utiliser ces exemples, créez un volume et utilisez les noms de catalogue, de schéma et de volume de ce volume pour définir le chemin d'accès au volume utilisé par les exemples. Les utilisateurs de l'édition gratuite ont accès au catalogue du Workspace et au schéma
defaultdefault. -
Vous devez disposer des autorisations suivantes dans Unity Catalog :
READ VOLUMEetWRITE VOLUMEpour le volume utilisé pour ce tutorielUSE SCHEMApour le schéma utilisé pour ce tutorielUSE CATALOGpour le catalogue utilisé pour ce tutoriel
Pour définir ces autorisations, consultez votre administrateur Databricks ou la référence des privilèges Unity Catalog. Les utilisateurs de l’édition gratuite disposent de ces privilèges sur le catalogue de Workspace et le schéma
defaultdefault.
Pour un Notebook complet sur cet article, consultez les Notebooks tutoriels sur le DataFrame.
Étape 1 : Définir des variables et charger le fichier CSV
Cette étape définit des variables à utiliser dans ce tutoriel, puis charge un fichier CSV contenant des données de prénoms de bébés provenant de health.data.ny.gov dans votre volume Unity Catalog. Vous avez besoin des noms d'un catalogue, d'un schéma et d'un volume Unity Catalog.
Si vous ne connaissez pas les noms de votre catalogue et de votre schéma, cliquez sur Catalogue dans le panneau latéral. Le catalogue de Workspace partage un nom avec votre Workspace et est listé dans le panneau latéral. Développez-le pour voir les schémas disponibles. Les utilisateurs de Free Edition et d'essai gratuit peuvent utiliser le catalogue Workspace et le schéma
default.
Si vous n’avez pas de volume, créez-en un en exécutant la commande suivante dans une cellule de notebook (remplacez <catalog_name> et <schema_name> par vos valeurs) :
CREATE VOLUME IF NOT EXISTS <catalog_name>.<schema_name>.my_volume
-
Ouvrez un nouveau Notebook en cliquant sur l'icône
. Pour apprendre à naviguer dans les Notebooks Databricks, voir Personnaliser l'apparence des Notebooks Databricks.
-
Copiez et collez le code suivant dans la nouvelle cellule de notebook vide. Remplacez
<catalog-name>,<schema-name>et<volume-name>par les noms de catalogue, de schéma et de volume pour un volume Unity Catalog. Remplacez<table_name>par un nom de table de votre choix. Vous chargerez des données de noms de bébé dans cette table plus tard dans ce tutoriel.
- Python
- Scala
- R
catalog = "<catalog_name>"
schema = "<schema_name>"
volume = "<volume_name>"
download_url = "https://health.data.ny.gov/api/views/jxy9-yhdk/rows.csv"
file_name = "rows.csv"
table_name = "<table_name>"
path_volume = "/Volumes/" + catalog + "/" + schema + "/" + volume
path_table = catalog + "." + schema
print(path_table) # Show the complete path
print(path_volume) # Show the complete path
val catalog = "<catalog_name>"
val schema = "<schema_name>"
val volume = "<volume_name>"
val downloadUrl = "https://health.data.ny.gov/api/views/jxy9-yhdk/rows.csv"
val fileName = "rows.csv"
val tableName = "<table_name>"
val pathVolume = s"/Volumes/$catalog/$schema/$volume"
val pathTable = s"$catalog.$schema"
print(pathVolume) // Show the complete path
print(pathTable) // Show the complete path
catalog <- "<catalog_name>"
schema <- "<schema_name>"
volume <- "<volume_name>"
download_url <- "https://health.data.ny.gov/api/views/jxy9-yhdk/rows.csv"
file_name <- "rows.csv"
table_name <- "<table_name>"
path_volume <- paste("/Volumes/", catalog, "/", schema, "/", volume, sep = "")
path_table <- paste(catalog, ".", schema, sep = "")
print(path_volume) # Show the complete path
print(path_table) # Show the complete path
-
Appuyez sur
Shift+Enterpour exécuter la cellule et créer une nouvelle cellule vide. -
Chargez le fichier CSV dans votre volume. Choisissez l'une des méthodes suivantes :
- Upload using the Workspace UI — Utilisez cette méthode si vous utilisez Databricks Free Edition, ou si le download de code dans l'option B échoue avec une erreur réseau. Free Edition et les autres environnements de compute Serverless restreignent l'accès Internet sortant, vous devez donc upload le fichier depuis votre machine locale.
- download using code — Utilisez cette méthode si votre environnement de compute a un accès Internet sortant.
Option A : upload à l'aide de l'interface utilisateur de Workspace
- Sur votre machine locale, ouvrez health.data.ny.gov/api/views/jxy9-yhdk/rows.csv dans votre navigateur. Le fichier download sur votre ordinateur sous la forme de
rows.csv, ce qui correspond à la variablefile_namedéfinie précédemment. - Retournez à votre Workspace Databricks. Dans la barre latérale, cliquez sur
Nouveau > Ajouter ou upload de données .
- Cliquez sur upload les fichiers vers un volume .
- Cliquez sur **Parcourir** et sélectionnez le
rows.csvfichier, ou glissez-déposez-le dans la zone d'upload. - Sous Volume de destination , sélectionnez le volume que vous avez spécifié ci-dessus.
- Une fois l'upload terminé, retournez à votre notebook et poursuivez à l'étape 2.
Pour plus de détails sur l'upload de fichiers, consultez Utiliser des fichiers dans les volumes Unity Catalog.
Option B : download à l'aide de code
Copiez et collez le code suivant dans la nouvelle cellule de notebook vide. Ce code copie le fichier
rows.csvde health.data.ny.gov dans votre volume Unity Catalog à l'aide de la commande Databricks dbutils. Appuyez surShift+Enterpour exécuter la cellule, puis passer à la cellule suivante.
- Python
- Scala
- R
dbutils.fs.cp(f"{download_url}", f"{path_volume}/{file_name}")
dbutils.fs.cp(downloadUrl, s"$pathVolume/$fileName")
dbutils.fs.cp(download_url, paste(path_volume, "/", file_name, sep = ""))
Étape 2 : créer un DataFrame
Cette étape crée un DataFrame nommé df1 avec des données de test, puis affiche son contenu.
- Copiez et collez le code suivant dans la nouvelle cellule de notebook vide. Ce code crée le DataFrame avec des données de test, puis affiche le contenu et le schéma du DataFrame.
- Python
- Scala
- R
data = [[2021, "test", "Albany", "M", 42]]
columns = ["Year", "First_Name", "County", "Sex", "Count"]
df1 = spark.createDataFrame(data, schema="Year int, First_Name STRING, County STRING, Sex STRING, Count int")
display(df1) # The display() method is specific to Databricks notebooks and provides a richer visualization.
# df1.show() The show() method is a part of the Apache Spark DataFrame API and provides basic visualization.
val data = Seq((2021, "test", "Albany", "M", 42))
val columns = Seq("Year", "First_Name", "County", "Sex", "Count")
val df1 = data.toDF(columns: _*)
display(df1) // The display() method is specific to Databricks notebooks and provides a richer visualization.
// df1.show() The show() method is a part of the Apache Spark DataFrame API and provides basic visualization.
# Load the SparkR package that is already preinstalled on the cluster.
library(SparkR)
data <- data.frame(
Year = as.integer(c(2021)),
First_Name = c("test"),
County = c("Albany"),
Sex = c("M"),
Count = as.integer(c(42))
)
df1 <- createDataFrame(data)
display(df1) # The display() method is specific to Databricks notebooks and provides a richer visualization.
# head(df1) The head() method is a part of the Apache SparkR DataFrame API and provides basic visualization.
- Appuyez sur
Shift+Enterpour exécuter la cellule, puis passez à la cellule suivante.
Étape 3 : Chargez les données dans un DataFrame à partir d'un fichier CSV
Cette étape crée un DataFrame nommé df_csv à partir du fichier CSV que vous avez précédemment chargé dans votre volume Unity Catalog. Voir spark.read.csv.
- Copiez et collez le code suivant dans la nouvelle cellule de notebook vide. Ce code charge les données de noms de bébé dans le DataFrame
df_csvà partir du fichier CSV, puis affiche le contenu du DataFrame.
- Python
- Scala
- R
df_csv = spark.read.csv(f"{path_volume}/{file_name}",
header=True,
inferSchema=True,
sep=",")
display(df_csv)
val dfCsv = spark.read
.option("header", "true")
.option("inferSchema", "true")
.option("delimiter", ",")
.csv(s"$pathVolume/$fileName")
display(dfCsv)
df_csv <- read.df(paste(path_volume, "/", file_name, sep=""),
source="csv",
header = TRUE,
inferSchema = TRUE,
delimiter = ",")
display(df_csv)
- Appuyez sur
Shift+Enterpour exécuter la cellule, puis passez à la cellule suivante.
Vous pouvez charger des données à partir de nombreux formats de fichier pris en charge.
Étape 4 : Visualisez et interagissez avec votre DataFrame
Consultez et interagissez avec vos DataFrames baby names à l’aide des méthodes suivantes.
Afficher le schéma du DataFrame
Découvrez comment afficher le schéma d'un DataFrame Apache Spark. Apache Spark utilise le terme schéma pour désigner les noms et les types de données des colonnes du DataFrame.
Databricks utilise également le terme « schéma » pour décrire un ensemble de tables enregistrées dans un catalogue.
- Copiez et collez le code suivant dans une cellule de Notebook vide. Ce code montre le schéma de vos DataFrames avec la méthode
.printSchema()pour afficher les schémas des deux DataFrames - pour préparer l'union des deux DataFrames.
- Python
- Scala
- R
df_csv.printSchema()
df1.printSchema()
dfCsv.printSchema()
df1.printSchema()
printSchema(df_csv)
printSchema(df1)
- Appuyez sur
Shift+Enterpour exécuter la cellule, puis passez à la cellule suivante.
Renommer une colonne du DataFrame
Découvrez comment renommer une colonne dans un DataFrame.
- Copiez et collez le code suivant dans une cellule de Notebook vide. Ce code renomme une colonne dans le DataFrame
df1_csvpour qu'elle corresponde à la colonne respective dans le DataFramedf1. Ce code utilise la méthode Apache SparkwithColumnRenamed().
- Python
- Scala
- R
df_csv = df_csv.withColumnRenamed("First Name", "First_Name")
df_csv.printSchema()
val dfCsvRenamed = dfCsv.withColumnRenamed("First Name", "First_Name")
// when modifying a DataFrame in Scala, you must assign it to a new variable
dfCsvRenamed.printSchema()
df_csv <- withColumnRenamed(df_csv, "First Name", "First_Name")
printSchema(df_csv)
- Appuyez sur
Shift+Enterpour exécuter la cellule, puis passez à la cellule suivante.
Combiner des DataFrames
Apprenez à créer un nouveau DataFrame qui ajoute les lignes d'un DataFrame à un autre.
- Copiez et collez le code suivant dans une cellule de Notebook vide. Ce code utilise la méthode
union()d’Apache Spark pour combiner le contenu de votre premier DataFramedfavec le DataFramedf_csvcontenant les données des prénoms chargées à partir du fichier CSV.
- Python
- Scala
- R
df = df1.union(df_csv)
display(df)
val df = df1.union(dfCsvRenamed)
display(df)
display(df <- union(df1, df_csv))
- Appuyez sur
Shift+Enterpour exécuter la cellule, puis passez à la cellule suivante.
Filtrer les lignes dans un DataFrame
Découvrez les prénoms de bébé les plus populaires dans votre ensemble de données en filtrant les lignes, en utilisant les méthodes Apache Spark .filter() ou .where(). Utilisez le filtrage pour sélectionner un sous-ensemble de lignes à retourner ou à modifier dans un DataFrame. Il n'y a pas de différence de performances ou de syntaxe, comme le montrent les exemples suivants.
Utilisation de .filter() méthode
- Copiez et collez le code suivant dans une cellule de Notebook vide. Ce code utilise la méthode
.filter()d'Apache Spark pour afficher les lignes dans le DataFrame avec un nombre supérieur à 50.
- Python
- Scala
- R
display(df.filter(df["Count"] > 50))
display(df.filter(df("Count") > 50))
display(filteredDF <- filter(df, df$Count > 50))
- Appuyez sur
Shift+Enterpour exécuter la cellule, puis passez à la cellule suivante.
Utilisation de .where() méthode
- Copiez et collez le code suivant dans une cellule de Notebook vide. Ce code utilise la méthode
.where()d'Apache Spark pour afficher les lignes dans le DataFrame avec un nombre supérieur à 50.
- Python
- Scala
- R
display(df.where(df["Count"] > 50))
display(df.where(df("Count") > 50))
display(filtered_df <- where(df, df$Count > 50))
- Appuyez sur
Shift+Enterpour exécuter la cellule, puis passez à la cellule suivante.
Sélectionnez des colonnes à partir d'un DataFrame et classez-les par fréquence
Découvrez la fréquence des prénoms avec la méthode select() afin de spécifier les colonnes du DataFrame à retourner. Utilisez les fonctions orderby et desc d'Apache Spark pour trier les résultats.
Le module pyspark.sql pour Apache Spark prend en charge les fonctions SQL. Parmi ces fonctions que nous utilisons dans ce didacticiel figurent les fonctions Apache Spark orderBy(), desc() et expr(). Vous activez l'utilisation de ces fonctions en les important dans votre session selon les besoins.
- Copiez et collez le code suivant dans une cellule de Notebook vide. Ce code importe la fonction
desc()puis utilise la méthodeselect()Apache Spark et les fonctionsorderBy()etdesc()Apache Spark pour afficher les noms les plus courants et leurs décomptes par ordre décroissant.
- Python
- Scala
- R
from pyspark.sql.functions import desc
display(df.select("First_Name", "Count").orderBy(desc("Count")))
import org.apache.spark.sql.functions.desc
display(df.select("First_Name", "Count").orderBy(desc("Count")))
display(arrange(select(df, df$First_Name, df$Count), desc(df$Count)))
- Appuyez sur
Shift+Enterpour exécuter la cellule, puis passez à la cellule suivante.
Créer un DataFrame de sous-ensemble
Apprenez à créer un sous-DataFrame à partir d’un DataFrame existant.
- Copiez et collez le code suivant dans une cellule de Notebook vide. Ce code utilise la méthode
filterd'Apache Spark pour créer un nouveau DataFrame restreignant les données par année, nombre et sexe. Elle utilise la méthode Apache Sparkselect()pour limiter les colonnes. Il utilise également les fonctionsorderBy()etdesc()d’Apache Spark pour trier le nouveau DataFrame par nombre.
- Python
- Scala
- R
subsetDF = df.filter((df["Year"] == 2009) & (df["Count"] > 100) & (df["Sex"] == "F")).select("First_Name", "County", "Count").orderBy(desc("Count"))
display(subsetDF)
val subsetDF = df.filter((df("Year") === 2009) && (df("Count") > 100) && (df("Sex") === "F")).select("First_Name", "County", "Count").orderBy(desc("Count"))
display(subsetDF)
subsetDF <- select(filter(df, (df$Count > 100) & (df$year == 2009) & df["Sex"] == "F")), "First_Name", "County", "Count")
display(subsetDF)
- Appuyez sur
Shift+Enterpour exécuter la cellule, puis passez à la cellule suivante.
Étape 5 : Enregistrer le DataFrame
Apprenez à enregistrer un DataFrame. Vous pouvez soit enregistrer votre DataFrame dans une table, soit écrire le DataFrame dans un fichier ou plusieurs fichiers.
Enregistrer le DataFrame dans une table
Databricks utilise le format Delta Lake pour toutes les tables par default. Pour enregistrer votre DataFrame, vous devez disposer des privilèges de table CREATE sur le catalogue et le schéma.
- Copiez et collez le code suivant dans une cellule de Notebook vide. Ce code enregistre le contenu du DataFrame dans une table en utilisant la variable que vous avez définie au start de ce tutoriel.
- Python
- Scala
- R
df.write.mode("overwrite").saveAsTable(f"{path_table}.{table_name}")
df.write.mode("overwrite").saveAsTable(s"$pathTable" + "." + s"$tableName")
saveAsTable(df, paste(path_table, ".", table_name), mode = "overwrite")
- Appuyez sur
Shift+Enterpour exécuter la cellule, puis passez à la cellule suivante.
La plupart des Spark applications Apache fonctionnent sur de grands jeux de données et de manière distribuée. Apache Spark écrit un répertoire de fichiers plutôt qu'un seul fichier. Delta Lake divise les dossiers et les fichiers Parquet. De nombreux systèmes de données peuvent lire ces répertoires de fichiers. Databricks recommande d'utiliser des tables plutôt que des chemins de fichiers pour la plupart des applications.
Enregistrer le DataFrame dans des fichiers JSON
- Copiez et collez le code suivant dans une cellule de Notebook vide. Ce code enregistre le DataFrame dans un répertoire de fichiers JSON.
- Python
- Scala
- R
df.write.format("json").mode("overwrite").save("/tmp/json_data")
df.write.format("json").mode("overwrite").save("/tmp/json_data")
write.df(df, path = "/tmp/json_data", source = "json", mode = "overwrite")
- Appuyez sur
Shift+Enterpour exécuter la cellule, puis passez à la cellule suivante.
Lire le DataFrame à partir d’un fichier JSON
Découvrez comment utiliser la méthode spark.read.format() d’Apache Spark pour lire les données JSON d’un répertoire dans un DataFrame.
- Copiez et collez le code suivant dans une cellule de Notebook vide. Ce code affiche les fichiers JSON que vous avez enregistrés dans l'exemple précédent.
- Python
- Scala
- R
display(spark.read.format("json").json("/tmp/json_data"))
display(spark.read.format("json").json("/tmp/json_data"))
display(read.json("/tmp/json_data"))
- Appuyez sur
Shift+Enterpour exécuter la cellule, puis passez à la cellule suivante.
Tâches supplémentaires : Exécutez des requêtes SQL dans PySpark, Scala et R
Les DataFrames Apache Spark offrent les options suivantes pour combiner SQL avec PySpark, Scala et R. Vous pouvez exécuter le code suivant dans le même Notebook que celui que vous avez créé pour ce tutoriel.
Spécifier une colonne comme query SQL
Découvrez comment utiliser la méthode selectExpr() d'Apache Spark. Ceci est une variante de la méthode select() qui accepte des expressions SQL et renvoie un DataFrame mis à jour. Cette méthode vous permet d'utiliser une expression SQL, telle que upper.
- Copiez et collez le code suivant dans une cellule de Notebook vide. Ce code utilise la méthode
selectExpr()d'Apache Spark et l'expressionupperSQL pour convertir une colonne de chaîne de caractères en majuscules (et renommer la colonne).
- Python
- Scala
- R
display(df.selectExpr("Count", "upper(County) as big_name"))
display(df.selectExpr("Count", "upper(County) as big_name"))
display(df_selected <- selectExpr(df, "Count", "upper(County) as big_name"))
- Appuyez sur
Shift+Enterpour exécuter la cellule, puis passez à la cellule suivante.
Utilisez expr() pour la syntaxe SQL pour une colonne
Découvrez comment importer et utiliser la fonction Apache Spark expr() pour utiliser la syntaxe SQL partout où une colonne serait spécifiée.
- Copiez et collez le code suivant dans une cellule de Notebook vide. Ce code importe la fonction
expr(), puis utilise la fonction Apache Sparkexpr()et l'expression SQLlowerpour convertir une colonne de chaîne en minuscules (et renommer la colonne).
- Python
- Scala
- R
from pyspark.sql.functions import expr
display(df.select("Count", expr("lower(County) as little_name")))
import org.apache.spark.sql.functions.{col, expr}
// Scala requires us to import the col() function as well as the expr() function
display(df.select(col("Count"), expr("lower(County) as little_name")))
display(df_selected <- selectExpr(df, "Count", "lower(County) as little_name"))
# expr() function is not supported in R, selectExpr in SparkR replicates this functionality
- Appuyez sur
Shift+Enterpour exécuter la cellule, puis passez à la cellule suivante.
Exécuter une requête SQL arbitraire à l’aide de spark.sql() fonction
Découvrez comment utiliser la fonction Apache Spark spark.sql() pour exécuter des requêtes SQL arbitraires.
- Copiez et collez le code suivant dans une cellule de Notebook vide. Ce code utilise la fonction Apache Spark
spark.sql()pour interroger une table SQL à l'aide de la syntaxe SQL.
- Python
- Scala
- R
display(spark.sql(f"SELECT * FROM {path_table}.{table_name}"))
display(spark.sql(s"SELECT * FROM $pathTable.$tableName"))
display(sql(paste("SELECT * FROM", path_table, ".", table_name)))
- Appuyez sur
Shift+Enterpour exécuter la cellule, puis passez à la cellule suivante.
Notebooks de tutoriel DataFrame
Les notebooks suivants incluent les exemples de query de ce tutoriel.
- Python
- Scala
- R