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 SCHEMAetCREATE TABLEsur le catalogueworkspace. 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.
- Visitez la page Enregistrements de personnes synthétiques : 10K à 10M enregistrements sur Kaggle.
- Cliquez sur Download , puis sur Download dataset au format zip . Cela télécharge un fichier nommé
archive.zipsur votre machine locale. - Extrayez le dossier
archivedu fichierarchive.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.
- Ouvrez l'Explorateur de catalogues en cliquant sur
Catalog dans la barre latérale.
- Dans l'Explorateur de catalogues, cliquez sur
Ajouter des données et Créer un volume .
- Nommez le volume
my-volumeet sélectionnez **Volume géré** comme type de volume. - Sélectionnez le catalogue et le
workspacedefaultschéma, puis cliquez sur **Créer**. - Ouvrez
my-volumeet cliquez sur upload vers ce volume . - Faites glisser et déposez ou parcourez pour sélectionner le fichier
person_10000.csvdans le dossierarchivede votre machine locale. - Click upload .
Enfin, créez un notebook pour l'exécution du code d'exemple.
- Cliquez sur
Nouveau dans la barre latérale.
- Cliquez sur
Notebook pour créer un nouveau Notebook.
- 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
- Scala
- SQL
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)
import org.apache.spark.sql.types._
val schema = StructType(Array(
StructField("id", IntegerType, true),
StructField("firstName", StringType, true),
StructField("lastName", StringType, true),
StructField("gender", StringType, true),
StructField("age", IntegerType, true)
))
val 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.saveAsTable("workspace.default.people_10k")
// View the new table.
val df2 = spark.read.table("workspace.default.people_10k")
display(df2)
-- Create the table with only the required columns and rename person_id to id.
CREATE OR REPLACE TABLE workspace.default.people_10k AS
SELECT
person_id AS id,
firstname,
lastname,
gender,
age
FROM read_files(
'/Volumes/workspace/default/my-volume/person_10000.csv',
format => 'csv',
header => true
);
-- View the new table.
SELECT * FROM workspace.default.people_10k;
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.
CREATE TABLE workspace.default.people_10k_prod LIKE workspace.default.people_10k
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
- Scala
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"))
import io.delta.tables.DeltaTable
DeltaTable.createOrReplace(spark)
.tableName("workspace.default.people_10k")
.addColumn("id", "INT")
.addColumn("firstName", "STRING")
.addColumn(
DeltaTable.columnBuilder("lastName")
.dataType("STRING")
.comment("surname")
.build()
)
.addColumn("gender", "STRING")
.addColumn("age", "INT")
.execute()
display(spark.read.table("workspace.default.people_10k"))
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
- Scala
- SQL
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)
import org.apache.spark.sql.types._
import io.delta.tables._
// Define schema
val schema = StructType(Array(
StructField("id", IntegerType, true),
StructField("firstName", StringType, true),
StructField("lastName", StringType, true),
StructField("gender", StringType, true),
StructField("age", IntegerType, true)
))
// Create data as Seq of Tuples
val data = Seq(
(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 DataFrame directly from Seq of Tuples
val people_10k_updates = spark.createDataFrame(data).toDF(
"id", "firstName", "lastName", "gender", "age"
)
people_10k_updates.createOrReplaceTempView("people_10k_updates")
// Merge the source and target tables
val deltaTable = DeltaTable.forName(spark, "workspace.default.people_10k")
deltaTable.as("people_10k")
.merge(
people_10k_updates.as("people_10k_updates"),
"people_10k.id = people_10k_updates.id"
)
.whenMatched()
.updateAll()
.whenNotMatched()
.insertAll()
.execute()
// View the additions to the table.
val df = spark.read.table("workspace.default.people_10k")
val df_filtered = df.filter($"id" >= 10001)
display(df_filtered)
-- Create the source table if it does not exist. Otherwise, replace the existing source table.
CREATE OR REPLACE TABLE workspace.default.people_10k_updates(
id INT,
firstName STRING,
lastName STRING,
gender STRING,
age INT
);
-- Insert new data into the source table.
INSERT INTO workspace.default.people_10k_updates VALUES
(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);
-- Merge the source and target tables.
MERGE INTO workspace.default.people_10k AS people_10k
USING workspace.default.people_10k_updates AS people_10k_updates
ON people_10k.id = people_10k_updates.id
WHEN MATCHED THEN
UPDATE SET *
WHEN NOT MATCHED THEN
INSERT *;
-- View the additions to the table.
SELECT * FROM workspace.default.people_10k WHERE id >= 10001
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
- Scala
- SQL
people_df = spark.read.table("workspace.default.people_10k")
display(people_df)
val people_df = spark.read.table("workspace.default.people_10k")
display(people_df)
SELECT * FROM workspace.default.people_10k;
É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
- Scala
- SQL
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)
// Create the new data.
val data = Seq(
(10007, "Miku", "Hatsune", "F", 25)
)
val df = spark.createDataFrame(data)
.toDF("id", "firstName", "lastName", "gender", "age")
// Append the new data to the target table
df.write.mode("append").saveAsTable("workspace.default.people_10k")
// View the new addition.
val df2 = spark.read.table("workspace.default.people_10k")
val df_filtered = df2.filter($"id" === 10007)
display(df_filtered)
CREATE OR REPLACE TABLE workspace.default.people_10k_new (
id INT,
firstName STRING,
lastName STRING,
gender STRING,
age INT
);
-- Insert the new data.
INSERT INTO workspace.default.people_10k_new VALUES
(10007, 'Miku', 'Hatsune', 'F', 25);
-- Append the new data to the target table.
INSERT INTO workspace.default.people_10k
SELECT * FROM workspace.default.people_10k_new;
-- View the new addition.
SELECT * FROM workspace.default.people_10k WHERE id = 10007;
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
- Scala
- SQL
df.write.mode("overwrite").saveAsTable("workspace.default.people_10k")
df.write.mode("overwrite").saveAsTable("workspace.default.people_10k")
INSERT OVERWRITE TABLE workspace.default.people_10k SELECT * FROM workspace.default.people_10k_2
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
- Scala
- SQL
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)
import io.delta.tables._
import org.apache.spark.sql.functions._
val deltaTable = DeltaTable.forName(spark, "workspace.default.people_10k")
// Declare the predicate and update rows using a SQL-formatted string.
deltaTable.updateExpr(
"gender = 'Female'",
Map("gender" -> "'F'")
)
// Declare the predicate and update rows using Spark SQL functions.
deltaTable.update(
col("gender") === "Male",
Map("gender" -> lit("M")));
deltaTable.update(
col("gender") === "Other",
Map("gender" -> lit("O")));
// View the updated table.
val df = spark.read.table("workspace.default.people_10k")
display(df)
-- Declare the predicate and update rows.
UPDATE workspace.default.people_10k SET gender = 'F' WHERE gender = 'Female';
UPDATE workspace.default.people_10k SET gender = 'M' WHERE gender = 'Male';
UPDATE workspace.default.people_10k SET gender = 'O' WHERE gender = 'Other';
-- View the updated table.
SELECT * FROM workspace.default.people_10k;
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
- Scala
- SQL
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)
import io.delta.tables._
import org.apache.spark.sql.functions._
val 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.
val df = spark.read.table("workspace.default.people_10k")
display(df)
-- Delete rows using a predicate.
DELETE FROM workspace.default.people_10k WHERE age < '21';
-- View the updated table.
SELECT * FROM workspace.default.people_10k;
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
- Scala
- SQL
from delta.tables import *
deltaTable = DeltaTable.forName(spark, "workspace.default.people_10k")
display(deltaTable.history())
import io.delta.tables._
val deltaTable = DeltaTable.forName(spark, "workspace.default.people_10k")
display(deltaTable.history())
DESCRIBE HISTORY workspace.default.people_10k
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
- Scala
- SQL
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'"))
import io.delta.tables._
val deltaTable = DeltaTable.forName(spark, "workspace.default.people_10k")
val 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'"))
-- Query using the version number
SELECT * FROM workspace.default.people_10k VERSION AS OF 0;
-- Query using the timestamp
SELECT * FROM workspace.default.people_10k TIMESTAMP AS OF '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
- Scala
- SQL
# 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)
// Query using the version number.
val dfVersion = spark.read
.option("versionAsOf", 0)
.table("workspace.default.people_10k")
// Query using the timestamp.
val dfTimestamp = spark.read
.option("timestampAsOf", "2026-01-05T23:09:47.000+00:00")
.table("workspace.default.people_10k")
display(dfVersion)
display(dfTimestamp)
-- Create a temporary view from version 0 of the table.
CREATE OR REPLACE TEMPORARY VIEW people_10k_v0 AS
SELECT * FROM workspace.default.people_10k VERSION AS OF 0;
-- Create a temporary view from a previous timestamp of the table.
CREATE OR REPLACE TEMPORARY VIEW people_10k_t0 AS
SELECT * FROM workspace.default.people_10k TIMESTAMP AS OF '2026-01-05T23:09:47.000+00:00';
SELECT * FROM people_10k_v0;
SELECT * FROM people_10k_t0;
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
- Scala
- SQL
from delta.tables import *
deltaTable = DeltaTable.forName(spark, "workspace.default.people_10k")
deltaTable.optimize().executeCompaction()
import io.delta.tables._
val deltaTable = DeltaTable.forName(spark, "workspace.default.people_10k")
deltaTable.optimize().executeCompaction()
OPTIMIZE workspace.default.people_10k
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
- Scala
- SQL
spark.sql("ALTER TABLE workspace.default.people_10k CLUSTER BY (firstName)")
spark.sql("OPTIMIZE FULL workspace.default.people_10k")
spark.sql("ALTER TABLE workspace.default.people_10k CLUSTER BY (firstName)")
spark.sql("OPTIMIZE FULL workspace.default.people_10k")
ALTER TABLE workspace.default.people_10k CLUSTER BY (firstName);
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
- Scala
- SQL
from delta.tables import *
deltaTable = DeltaTable.forName(spark, "workspace.default.people_10k")
deltaTable.vacuum()
import io.delta.tables._
val deltaTable = DeltaTable.forName(spark, "workspace.default.people_10k")
deltaTable.vacuum()
VACUUM workspace.default.people_10k
Pour plus d'informations sur l'utilisation efficace de l'opération vacuum, voir Supprimer les fichiers de données inutilisés avec vacuum.