Aller au contenu principal

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.

remarque

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 :

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 default default.

  • Vous devez disposer des autorisations suivantes dans Unity Catalog :

    • READ VOLUME et WRITE VOLUME pour le volume utilisé pour ce tutoriel
    • USE SCHEMA pour le schéma utilisé pour ce tutoriel
    • USE CATALOG pour 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 default default.

astuce

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.

astuce

Si vous ne connaissez pas les noms de votre catalogue et de votre schéma, cliquez sur Icône de données. 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) :

SQL
CREATE VOLUME IF NOT EXISTS <catalog_name>.<schema_name>.my_volume
  1. Ouvrez un nouveau Notebook en cliquant sur l'icône Nouvelle icône. Pour apprendre à naviguer dans les Notebooks Databricks, voir Personnaliser l'apparence des Notebooks Databricks.

  2. 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
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
  1. Appuyez sur Shift+Enter pour exécuter la cellule et créer une nouvelle cellule vide.

  2. 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

    1. 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 variable file_name définie précédemment.
    2. Retournez à votre Workspace Databricks. Dans la barre latérale, cliquez sur Nouvelle icône Nouveau > Ajouter ou upload de données .
    3. Cliquez sur upload les fichiers vers un volume .
    4. Cliquez sur **Parcourir** et sélectionnez le rows.csv fichier, ou glissez-déposez-le dans la zone d'upload.
    5. Sous Volume de destination , sélectionnez le volume que vous avez spécifié ci-dessus.
    6. 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.csv de health.data.ny.gov dans votre volume Unity Catalog à l'aide de la commande Databricks dbutils. Appuyez sur Shift+Enter pour exécuter la cellule, puis passer à la cellule suivante.

Python
dbutils.fs.cp(f"{download_url}", f"{path_volume}/{file_name}")

Étape 2 : créer un DataFrame

Cette étape crée un DataFrame nommé df1 avec des données de test, puis affiche son contenu.

  1. 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
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.
  1. Appuyez sur Shift+Enter pour 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.

  1. 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
df_csv = spark.read.csv(f"{path_volume}/{file_name}",
header=True,
inferSchema=True,
sep=",")
display(df_csv)
  1. Appuyez sur Shift+Enter pour 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.

remarque

Databricks utilise également le terme « schéma » pour décrire un ensemble de tables enregistrées dans un catalogue.

  1. 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
df_csv.printSchema()
df1.printSchema()
  1. Appuyez sur Shift+Enter pour exécuter la cellule, puis passez à la cellule suivante.

Renommer une colonne du DataFrame

Découvrez comment renommer une colonne dans un DataFrame.

  1. Copiez et collez le code suivant dans une cellule de Notebook vide. Ce code renomme une colonne dans le DataFrame df1_csv pour qu'elle corresponde à la colonne respective dans le DataFrame df1. Ce code utilise la méthode Apache Spark withColumnRenamed().
Python
df_csv = df_csv.withColumnRenamed("First Name", "First_Name")
df_csv.printSchema()
  1. Appuyez sur Shift+Enter pour 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.

  1. 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 DataFrame df avec le DataFrame df_csv contenant les données des prénoms chargées à partir du fichier CSV.
Python
df = df1.union(df_csv)
display(df)
  1. Appuyez sur Shift+Enter pour 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

  1. 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
display(df.filter(df["Count"] > 50))
  1. Appuyez sur Shift+Enter pour exécuter la cellule, puis passez à la cellule suivante.

Utilisation de .where() méthode

  1. 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
display(df.where(df["Count"] > 50))
  1. Appuyez sur Shift+Enter pour 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.

  1. Copiez et collez le code suivant dans une cellule de Notebook vide. Ce code importe la fonction desc() puis utilise la méthode select() Apache Spark et les fonctions orderBy() et desc() Apache Spark pour afficher les noms les plus courants et leurs décomptes par ordre décroissant.
Python
from pyspark.sql.functions import desc
display(df.select("First_Name", "Count").orderBy(desc("Count")))
  1. Appuyez sur Shift+Enter pour 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.

  1. Copiez et collez le code suivant dans une cellule de Notebook vide. Ce code utilise la méthode filter d'Apache Spark pour créer un nouveau DataFrame restreignant les données par année, nombre et sexe. Elle utilise la méthode Apache Spark select() pour limiter les colonnes. Il utilise également les fonctions orderBy() et desc() d’Apache Spark pour trier le nouveau DataFrame par nombre.
Python
subsetDF = df.filter((df["Year"] == 2009) & (df["Count"] > 100) & (df["Sex"] == "F")).select("First_Name", "County", "Count").orderBy(desc("Count"))
display(subsetDF)
  1. Appuyez sur Shift+Enter pour 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.

  1. 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
df.write.mode("overwrite").saveAsTable(f"{path_table}.{table_name}")
  1. Appuyez sur Shift+Enter pour 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

  1. 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
df.write.format("json").mode("overwrite").save("/tmp/json_data")
  1. Appuyez sur Shift+Enter pour 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.

  1. 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
display(spark.read.format("json").json("/tmp/json_data"))
  1. Appuyez sur Shift+Enter pour 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.

  1. Copiez et collez le code suivant dans une cellule de Notebook vide. Ce code utilise la méthode selectExpr() d'Apache Spark et l'expression upper SQL pour convertir une colonne de chaîne de caractères en majuscules (et renommer la colonne).
Python
display(df.selectExpr("Count", "upper(County) as big_name"))
  1. Appuyez sur Shift+Enter pour 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.

  1. Copiez et collez le code suivant dans une cellule de Notebook vide. Ce code importe la fonction expr(), puis utilise la fonction Apache Spark expr() et l'expression SQL lower pour convertir une colonne de chaîne en minuscules (et renommer la colonne).
Python
from pyspark.sql.functions import expr
display(df.select("Count", expr("lower(County) as little_name")))
  1. Appuyez sur Shift+Enter pour 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.

  1. 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
display(spark.sql(f"SELECT * FROM {path_table}.{table_name}"))
  1. Appuyez sur Shift+Enter pour exécuter la cellule, puis passez à la cellule suivante.

Notebooks de tutoriel DataFrame

Les notebooks suivants incluent les exemples de query de ce tutoriel.

Ressources supplémentaires