Aller au contenu principal

Tutoriel : créer et gérer des tables Delta Lake

Ce tutoriel présente les opérations courantes sur les tables Delta Lake à l’aide de données d’exemple. Delta Lake est la couche de stockage optimisée qui constitue la base des tables sur Databricks. Sauf indication contraire, toutes les tables sur Databricks sont des tables Delta Lake.

Avant de commencer

Pour terminer ce tutoriel, vous avez besoin :

  • Autorisation d'utiliser une ressource de compute existante ou de créer une nouvelle ressource de compute. See compute.
  • Autorisations Unity Catalog : USE CATALOG, USE SCHEMA et CREATE TABLE sur le catalogue workspace. Pour définir ces autorisations, consultez votre administrateur Databricks ou la référence des privilèges Unity Catalog.

Ces exemples s’appuient sur un dataset appelé Synthetic Person Records: 10K to 10M Records . Ce dataset contient des enregistrements fictifs de personnes, notamment leur prénom et nom, leur sexe et leur âge.

Tout d'abord, download le dataset pour ce tutoriel.

  1. Visitez la page Enregistrements de personnes synthétiques : 10K à 10M enregistrements sur Kaggle.
  2. Cliquez sur Download , puis sur Download dataset au format zip . Cela télécharge un fichier nommé archive.zip sur votre machine locale.
  3. Extrayez le dossier archive du fichier archive.zip.

Ensuite, upload le dataset person_10000.csv vers un volume Unity Catalog dans votre Workspace Databricks. Databricks recommande d'upload vos données vers un volume Unity Catalog, car les volumes offrent des fonctionnalités pour accéder, stocker, régir et organiser les fichiers.

  1. Ouvrez l'Explorateur de catalogues en cliquant sur Icône de données. Catalog dans la barre latérale.
  2. Dans l'Explorateur de catalogues, cliquez sur Icône Ajouter ou plus Ajouter des données et Créer un volume .
  3. Nommez le volume my-volume et sélectionnez **Volume géré** comme type de volume.
  4. Sélectionnez le catalogue et le workspace default schéma, puis cliquez sur **Créer**.
  5. Ouvrez my-volume et cliquez sur upload vers ce volume .
  6. Faites glisser et déposez ou parcourez pour sélectionner le fichier person_10000.csv dans le dossier archive de votre machine locale.
  7. Click upload .

Enfin, créez un notebook pour l'exécution du code d'exemple.

  1. Cliquez sur Icône Ajouter ou plus Nouveau dans la barre latérale.
  2. Cliquez sur Icône du Notebook. Notebook pour créer un nouveau Notebook.
  3. Choisissez une langue pour le Notebook.

Créer une table

Créez une nouvelle table gérée Unity Catalog nommée workspace.default.people_10k à partir de person_10000.csv. Delta Lake est le default pour toutes les commandes de création, de lecture et d'écriture de tables dans Databricks.

Python
from pyspark.sql.types import StructType, StructField, IntegerType, StringType

schema = StructType([
StructField("id", IntegerType(), True),
StructField("firstName", StringType(), True),
StructField("lastName", StringType(), True),
StructField("gender", StringType(), True),
StructField("age", IntegerType(), True)
])

df = spark.read.format("csv").option("header", True).schema(schema).load("/Volumes/workspace/default/my-volume/person_10000.csv")

# Create the table if it does not exist. Otherwise, replace the existing table.
df.writeTo("workspace.default.people_10k").createOrReplace()

# If you know the table does not already exist, you can use this command instead.
# df.write.saveAsTable("workspace.default.people_10k")

# View the new table.
df = spark.read.table("workspace.default.people_10k")
display(df)

Il existe plusieurs façons de créer ou de cloner des tables. Pour plus d’informations, consultez CREATE TABLE.

Dans Databricks Runtime 13.3 LTS et versions ultérieures, vous pouvez utiliser CREATE TABLE LIKE pour créer une nouvelle table Delta Lake vide qui duplique le schéma et les propriétés de table d'une table Delta Lake source. Ceci peut être utile lors de la promotion de tables d'un environnement de développement vers la production.

SQL
CREATE TABLE workspace.default.people_10k_prod LIKE workspace.default.people_10k
info

Aperçu

Cette fonctionnalité est en aperçu public.

Utilisez l'DeltaTableBuilder API pour Python et Scala afin de créer une table vide. Comparée à DataFrameWriter et DataFrameWriterV2, l'API DeltaTableBuilder facilite la spécification d'informations supplémentaires comme les commentaires de colonne, les propriétés de table et les colonnes générées.

Python
from delta.tables import DeltaTable

(
DeltaTable.createIfNotExists(spark)
.tableName("workspace.default.people_10k_2")
.addColumn("id", "INT")
.addColumn("firstName", "STRING")
.addColumn("lastName", "STRING", comment="surname")
.addColumn("gender", "STRING")
.addColumn("age", "INT")
.execute()
)

display(spark.read.table("workspace.default.people_10k_2"))

Mettre à jour/insérer dans une table

Modifiez les enregistrements existants dans une table ou ajoutez-en de nouveaux à l’aide d’une opération appelée upsert . Pour merge un ensemble de mises à jour et d’insertions dans une table Delta Lake existante, utilisez la méthode DeltaTable.merge en Python et en Scala et l’instruction MERGE INTO en SQL.

Par exemple, Merge les données de la table source people_10k_updates vers la table Delta Lake cible workspace.default.people_10k. Lorsqu'il y a une ligne correspondante dans les deux tables, Delta Lake met à jour la colonne de données en utilisant l'expression donnée. Lorsqu'il n'y a pas de ligne correspondante, Delta Lake ajoute une nouvelle ligne.

Python
from pyspark.sql.types import StructType, StructField, StringType, IntegerType
from delta.tables import DeltaTable

schema = StructType([
StructField("id", IntegerType(), True),
StructField("firstName", StringType(), True),
StructField("lastName", StringType(), True),
StructField("gender", StringType(), True),
StructField("age", IntegerType(), True)
])

data = [
(10001, 'Billy', 'Luppitt', 'M', 55),
(10002, 'Mary', 'Smith', 'F', 98),
(10003, 'Elias', 'Leadbetter', 'M', 48),
(10004, 'Jane', 'Doe', 'F', 30),
(10005, 'Joshua', '', 'M', 90),
(10006, 'Ginger', '', 'F', 16),
]

# Create the source table if it does not exist. Otherwise, replace the existing source table.
people_10k_updates = spark.createDataFrame(data, schema)
people_10k_updates.createOrReplaceTempView("people_10k_updates")

# Merge the source and target tables.
deltaTable = DeltaTable.forName(spark, 'workspace.default.people_10k')

(deltaTable.alias("people_10k")
.merge(
people_10k_updates.alias("people_10k_updates"),
"people_10k.id = people_10k_updates.id")
.whenMatchedUpdateAll()
.whenNotMatchedInsertAll()
.execute()
)

# View the additions to the table.
df = spark.read.table("workspace.default.people_10k")
df_filtered = df.filter(df["id"] >= 10001)
display(df_filtered)

En SQL, l'opérateur * met à jour ou insère toutes les colonnes dans la table cible, en supposant que la table source a les mêmes colonnes que la table cible. Si la table cible n'a pas les mêmes colonnes, la query génère une erreur d'analyse. Vous devez également spécifier une valeur pour chaque colonne de votre table lorsque vous effectuez une opération d'insertion. Les valeurs de colonne peuvent être vides, par exemple, ''. Lorsque vous effectuez une opération d'insertion, vous n'avez pas besoin de mettre à jour toutes les valeurs.

Lire une table

Utilisez le nom de la table ou le chemin d'accès pour accéder aux données dans les tables Delta Lake. Pour accéder aux tables gérées par Unity Catalog, utilisez un nom de table entièrement qualifié. L’accès basé sur le chemin n’est pris en charge que pour les volumes et les tables externes, et non pour les tables gérées. Pour plus d’informations, consultez Règles de chemin d’accès et accès dans les volumes Unity Catalog.

Python
people_df = spark.read.table("workspace.default.people_10k")
display(people_df)

Écrire sur une table

Delta Lake utilise la syntaxe standard pour l'écriture de données dans les tables. Pour ajouter de nouvelles données à une table Delta Lake existante, utilisez le mode d'ajout. Contrairement à l'upsertion, l'écriture dans une table ne vérifie pas les enregistrements en double.

Python
from pyspark.sql.types import StructType, StructField, StringType, IntegerType
from pyspark.sql.functions import col

schema = StructType([
StructField("id", IntegerType(), True),
StructField("firstName", StringType(), True),
StructField("lastName", StringType(), True),
StructField("gender", StringType(), True),
StructField("age", IntegerType(), True)
])

data = [
(10007, 'Miku', 'Hatsune', 'F', 25)
]

# Create the new data.
df = spark.createDataFrame(data, schema)

# Append the new data to the target table.
df.write.mode("append").saveAsTable("workspace.default.people_10k")

# View the new addition.
df = spark.read.table("workspace.default.people_10k")
df_filtered = df.filter(df["id"] == 10007)
display(df_filtered)

Les sorties de cellule du Notebook Databricks affichent un maximum de 10 000 lignes ou de 2 Mo, la valeur la plus basse étant retenue. Puisque workspace.default.people_10k contient plus de 10 000 lignes, seules les 10 000 premières lignes apparaissent dans la sortie du Notebook pour display(df). Les lignes supplémentaires sont présentes dans la table, mais ne sont pas affichées dans la sortie du notebook en raison de cette limite. Vous pouvez afficher les lignes supplémentaires en les filtrant spécifiquement.

Pour remplacer toutes les données d'une table, utilisez le mode de remplacement.

Python
df.write.mode("overwrite").saveAsTable("workspace.default.people_10k")

Mettre à jour une table

Mettez à jour les données dans une table Delta Lake en fonction d'un prédicat. Par exemple, modifiez les valeurs de la colonne gender de Female à F, de Male à M et de Other à O.

Python
from delta.tables import *
from pyspark.sql.functions import *

deltaTable = DeltaTable.forName(spark, "workspace.default.people_10k")

# Declare the predicate and update rows using a SQL-formatted string.
deltaTable.update(
condition = "gender = 'Female'",
set = { "gender": "'F'" }
)

# Declare the predicate and update rows using Spark SQL functions.
deltaTable.update(
condition = col('gender') == 'Male',
set = { 'gender': lit('M') }
)

deltaTable.update(
condition = col('gender') == 'Other',
set = { 'gender': lit('O') }
)

# View the updated table.
df = spark.read.table("workspace.default.people_10k")
display(df)

Supprimer d'une table

Supprimez les données qui correspondent à un prédicat d'une table Delta Lake. Par exemple, le code ci-dessous illustre deux Opérations de suppression : la première supprime les lignes où l'âge est inférieur à 18, puis la seconde supprime les lignes où l'âge est inférieur à 21.

Python
from delta.tables import *
from pyspark.sql.functions import *

deltaTable = DeltaTable.forName(spark, "workspace.default.people_10k")

# Declare the predicate and delete rows using a SQL-formatted string.
deltaTable.delete("age < '18'")

# Declare the predicate and delete rows using Spark SQL functions.
deltaTable.delete(col('age') < '21')

# View the updated table.
df = spark.read.table("workspace.default.people_10k")
display(df)
important

La suppression retire les données de la dernière version de la table Delta Lake, mais ne les retire pas du stockage physique tant que les anciennes versions n'ont pas été explicitement vacuum. Pour plus d'informations, consultez vacuum.

Afficher l'historique de la table

Utilisez la méthode DeltaTable.history dans Python et Scala et l'instruction DESCRIBE HISTORY dans SQL pour afficher les informations de provenance de chaque écriture dans une table.

Python
from delta.tables import *

deltaTable = DeltaTable.forName(spark, "workspace.default.people_10k")
display(deltaTable.history())

query une version antérieure de la table à l'aide de time travel

Query un ancien cliché d'une table Delta Lake en utilisant le time travel Delta Lake. Pour interroger une version spécifique, utilisez le numéro de version de la table ou son Timestamp. Par exemple, query version 0 ou Timestamp 2026-01-05T23:09:47.000+00:00 de l'historique de la table.

Python
from delta.tables import *

deltaTable = DeltaTable.forName(spark, "workspace.default.people_10k")
deltaHistory = deltaTable.history()

# Query using the version number.
display(deltaHistory.where("version == 0"))

# Query using the timestamp.
display(deltaHistory.where("timestamp == '2026-01-05T23:09:47.000+00:00'"))

Pour les Timestamps, seules les chaînes de date ou de Timestamp sont acceptées. Par exemple, les chaînes de caractères doivent être formatées en "2026-01-05T22:43:15.000+00:00" ou "2026-01-05 22:43:15".

Utilisez les options DataFrameReader pour créer un DataFrame à partir d'une table Delta Lake qui est fixé à une version ou un timestamp spécifique de la table.

Python
# Query using the version number.
df = spark.read.option('versionAsOf', 0).table("workspace.default.people_10k")

# Query using the timestamp.
df = spark.read.option('timestampAsOf', '2026-01-05T23:09:47.000+00:00').table("workspace.default.people_10k")

display(df)

Pour plus d'informations, consultez Utiliser l'historique des tables.

Optimiser une table

Plusieurs modifications apportées à une table peuvent créer plusieurs petits fichiers, ce qui ralentit les performances des query de lecture. Utilisez l'opération d'optimisation pour améliorer la vitesse en combinant les petits fichiers en des fichiers plus grands. See OPTIMIZE.

Python
from delta.tables import *

deltaTable = DeltaTable.forName(spark, "workspace.default.people_10k")
deltaTable.optimize().executeCompaction()
remarque

Si l'optimisation prédictive est activée, vous n'avez pas besoin d'optimiser manuellement. L'optimisation prédictive gère automatiquement les tâches de maintenance. Pour plus d'informations, consultez l'optimisation prédictive pour les tables gérées par Unity Catalog.

Utilisez le Liquid Clustering

Pour améliorer davantage les performances de lecture, utilisez le clustering liquide pour colocaliser les données associées. Par exemple, activez le clustering sur la colonne à cardinalité élevée firstName. Utilisez OPTIMIZE FULL pour appliquer le clustering à toutes les données existantes. Pour plus d'informations, consultez Utiliser le clustering liquide pour les tables.

Python
spark.sql("ALTER TABLE workspace.default.people_10k CLUSTER BY (firstName)")
spark.sql("OPTIMIZE FULL workspace.default.people_10k")

Nettoyer les instantanés avec l'opération vacuum

Delta Lake dispose de l'isolation des instantanés pour les lectures, ce qui signifie qu'il est sans risque d'exécuter une opération d'optimisation pendant que d'autres utilisateurs ou jobs interrogent la table. Cependant, vous devriez nettoyer les anciens snapshots, car cela réduit les coûts de stockage, améliore les performances des query et assure la conformité des données. Exécutez l’opération VACUUM pour nettoyer les anciens snapshots. See vacuum.

Python
from delta.tables import *

deltaTable = DeltaTable.forName(spark, "workspace.default.people_10k")
deltaTable.vacuum()

Pour plus d'informations sur l'utilisation efficace de l'opération vacuum, voir Supprimer les fichiers de données inutilisés avec vacuum.

Ressources supplémentaires