Gérer l'historique des tables
Pour les tables Apache Iceberg et Delta Lake, chaque opération qui modifie une table crée une nouvelle version de table. Utilisez les informations d'historique pour auditer les opérations, restaurer une table ou query une table à un point précis dans le temps en utilisant le time travel.
N'utilisez pas l'historique de table comme solution de sauvegarde à long terme pour l'archivage des données. Utilisez uniquement les 7 derniers jours pour les opérations de time travel, sauf si vous avez défini des configurations de rétention des données et des logs sur une valeur plus élevée.
Récupérer l'historique de la table
Exécutez la commande DESCRIBE HISTORY pour récupérer des informations, notamment les opérations, l’utilisateur et le Timestamp de chaque écriture dans une table. Les opérations sont renvoyées dans l'ordre chronologique inverse.
Pour les colonnes renvoyées par DESCRIBE HISTORY, les valeurs de la colonne operationParameters et les métriques par opération dans la colonne operationMetrics, consultez Schéma de l’historique des tables et métriques d’opérations.
La rétention de l'historique des tables est déterminée par le paramètre de table logRetentionDuration, qui est 30 jours by default.
Le time travel et l'historique des tables sont régis par différents seuils de rétention. See time travel.
DESCRIBE HISTORY table_name -- get the full history of the table
DESCRIBE HISTORY table_name LIMIT 1 -- get the last operation only
Pour plus de détails sur la syntaxe Spark SQL, consultez DESCRIBE HISTORY.
Pour plus de détails sur la syntaxe Scala, Java et Python, consultez la documentation de l’API Delta Lake.
Catalog Explorer shows table history visually on the History tab.
Identifier le type d'OPTIMIZE opération
Le compactage automatique, le clustering fluide et le Z-ordering apparaissent tous dans l'historique de la table comme des opérations OPTIMIZE. Pour déterminer celle qui s'est exécutée, inspectez la colonne operationParameters.
Pour classifier chaque opération OPTIMIZE dans l'historique d'une table, exécutez ce qui suit :
SELECT
version,
timestamp,
CASE
WHEN operationParameters.clusterBy IS NOT NULL AND operationParameters.clusterBy <> '[]' THEN 'Liquid clustering'
WHEN operationParameters.zOrderBy IS NOT NULL AND operationParameters.zOrderBy <> '[]' THEN 'Z-ordering'
WHEN operationParameters.auto = 'true' THEN 'Auto compaction'
ELSE 'Manual OPTIMIZE'
END AS optimize_type,
operationParameters.auto AS is_auto_compaction,
operationParameters.clusterBy AS cluster_by,
operationParameters.zOrderBy AS z_order_by,
operationMetrics.numRemovedFiles AS files_compacted,
operationMetrics.numAddedFiles AS files_added,
operationMetrics.numRemovedBytes AS bytes_removed,
operationMetrics.numAddedBytes AS bytes_added
FROM (DESCRIBE HISTORY table_name)
WHERE operation = 'OPTIMIZE'
ORDER BY version DESC;
Les sections suivantes décrivent chaque valeur operationParameters en détail. Pour les définitions des clés operationMetrics que la query précédente sélectionne, consultez Opérations métriques.
Compactage automatique
Le compactage automatique définit le paramètre auto à true. Databricks Trigger automatiquement le compactage automatique après une écriture. Lorsque auto est false, un utilisateur ou un Job planifié a exécuté la commande OPTIMIZE.
Par exemple, une opération de compactage automatique affiche ce qui suit :
operationParameters: {
"auto": "true"
}
Pour plus d'informations sur le compactage automatique, consultez Compactage automatique.
Clustering fluide
Le clustering fluide remplit le parameter clusterBy avec les noms de colonne du clustering. Un tableau clusterBy vide ([]) indique uniquement le compactage de fichier.
Par exemple, une opération qui a regroupé les données par les colonnes date et region affiche ce qui suit :
operationParameters: {
"clusterBy": "[\"date\",\"region\"]"
}
Pour plus d'informations sur le clustering fluide, consultez Utiliser le clustering fluide pour les tables.
Z-ordering
Le Z-ordering remplit le paramètre zOrderBy avec les noms de colonne de l'ordre Z. Un tableau zOrderBy vide ([]) indique que l'opération n'a pas appliqué le Z-ordering.
Par exemple, une opération qui a appliqué le Z-ordering sur la colonne date affiche ce qui suit :
operationParameters: {
"zOrderBy": "[\"date\"]"
}
Portée de l'opération
Le predicate parameter indique si l'opération a été exécutée sur la table complète ou seulement sur une partie de celle-ci :
- Un tableau
predicatevide ([]) signifie que l'opération s'est exécutée sur la table entière. - Un tableau
predicaterenseigné signifie qu'une commandeOPTIMIZE table_name WHERE <partition_predicate>ciblée s'est exécutée uniquement sur les partitions qui correspondent au prédicat.
Par exemple, une opération ciblant les partitions qui correspondent à year = 2024 affiche ce qui suit :
operationParameters: {
"predicate": "[\"'year = 2024\"]"
}
time travel
Le voyage dans le temps prend en charge l’interrogation des versions précédentes des tables en fonction de l’horodatage ou de la version de la table (tel qu’enregistré dans le journal des transactions). Vous pouvez utiliser le time travel pour des applications telles que les suivantes :
- Recréer des analyses, des rapports ou des résultats, tels que le résultat d'un Modèle de machine learning. Ceci pourrait être utile pour le debugging ou l'audit, en particulier dans les secteurs d'activité réglementés.
- Rédaction de queries temporelles complexes.
- Correction des erreurs dans vos données.
- Fournir l'isolation d'instantanés pour un ensemble de query pour les tables à modification rapide.
Dans Databricks Runtime 18.0 et versions supérieures, les requêtes time travel sont bloquées si elles demandent une version antérieure à la propriété de table deletedFileRetentionDuration (7 jours par default). Pour les tables gérées par Unity Catalog, cela s'applique à Databricks Runtime 12.2 et versions ultérieures.
Syntaxe Time travel
Vous query une table avec time travel en ajoutant une clause après la spécification du nom de la table.
-
timestamp_expressionpeut être l'un des suivants :'2018-10-18T22:15:12.013Z', c'est-à-dire une chaîne qui peut être convertie en Timestampcast('2018-10-18 13:36:32 CEST' as timestamp)'2018-10-18', c'est-à-dire, une chaîne de datescurrent_timestamp() - interval 12 hoursdate_sub(current_date(), 1)- Toute autre expression qui est ou peut être castée en Timestamp
-
versionest une valeur longue qui peut être obtenue à partir de la sortie deDESCRIBE HISTORY table_spec.
Ni timestamp_expression ni version ne peuvent être des sous-requêtes.
Seules les chaînes de date ou de Timestamp sont acceptées. Par exemple, "2019-01-01" et "2019-01-01T00:00:00.000Z". Consultez le code suivant pour un exemple de syntaxe :
- SQL
- Python
SELECT * FROM people10m TIMESTAMP AS OF '2018-10-18T22:15:12.013Z';
SELECT * FROM people10m VERSION AS OF 123;
df1 = spark.read.option("timestampAsOf", "2019-01-01").table("people10m")
df2 = spark.read.option("versionAsOf", 123).table("people10m")
Vous pouvez également utiliser la syntaxe @ pour spécifier le Timestamp ou la version dans le nom de la table. Le Timestamp doit être au format yyyyMMddHHmmssSSS. Vous pouvez spécifier une version avec @v. Consultez le code suivant pour un exemple de syntaxe :
- SQL
- Python
-- Timestamp version
SELECT * FROM people10m@20190101000000000
-- Version number
SELECT * FROM people10m@v123
# Timestamp version
spark.read.table("people10m@20190101000000000")
# Version number
spark.read.table("people10m@v123")
Configurez la conservation des données pour les queries en time travel
Pour interroger une version précédente de table, vous devez conserver à la fois les logs et les fichiers de données pour cette version :
- Les fichiers de données sont supprimés lorsque
VACUUMs'exécute sur une table. - Les fichiers logs sont automatiquement supprimés après la mise en œuvre de points de contrôle des versions de table.
Pour augmenter le threshold de rétention des données pour les tables, vous devez configurer les propriétés de table suivantes, en remplaçant <format> par delta ou iceberg:
-
<format>.logRetentionDuration = "interval <interval>": contrôle la durée de conservation de l'historique d'une table. La default estinterval 30 days.- Dans Databricks Runtime 18.0 et versions ultérieures,
logRetentionDurationdoit être supérieur ou égal àdeletedFileRetentionDuration. Pour les tables gérées par Unity Catalog, cela s'applique à Databricks Runtime 12.2 et versions ultérieures.
- Dans Databricks Runtime 18.0 et versions ultérieures,
-
<format>.deletedFileRetentionDuration = "interval <interval>": détermine le threshold queVACUUMutilise pour supprimer les fichiers de données qui ne sont plus référencés dans la version actuelle de la table. The default estinterval 7 days.
Par exemple, pour accéder à 30 jours de données historiques, définissez delta.deletedFileRetentionDuration = "interval 30 days", ce qui correspond au default pour delta.logRetentionDuration.
L'augmentation du threshold de rétention des données peut entraîner une augmentation de vos coûts de stockage, car davantage de fichiers de données sont conservés.
Vous pouvez spécifier les propriétés de la table lors de la création de la table ou les définir avec une instruction ALTER TABLE. Consultez la référence des propriétés de table.
Exemples de time travel
Pour corriger les suppressions accidentelles d’une table pour l’utilisateur 111:
INSERT INTO my_table
SELECT * FROM my_table TIMESTAMP AS OF date_sub(current_date(), 1)
WHERE userId = 111
Pour corriger les mises à jour incorrectes accidentelles d'une table :
MERGE INTO my_table target
USING my_table TIMESTAMP AS OF date_sub(current_date(), 1) source
ON source.userId = target.userId
WHEN MATCHED THEN UPDATE SET *
Pour query le nombre de nouveaux clients ajoutés au cours de la dernière semaine :
SELECT
(
SELECT count(distinct userId)
FROM my_table
)
-
(
SELECT count(distinct userId)
FROM my_table TIMESTAMP AS OF date_sub(current_date(), 7)
) AS new_customers
Points de contrôle du log des transactions
Le journal des transactions enregistre les versions de table en tant que fichiers JSON dans le répertoire du journal des transactions, à côté des données de la table.
Pour optimiser l'interrogation des points de contrôle, les versions de table sont agrégées dans des fichiers Parquet de point de contrôle, ce qui améliore les performances en évitant la nécessité de lire toutes les versions JSON de l'historique de la table. Les utilisateurs n'ont pas besoin d'interagir directement avec les points de contrôle.
Databricks optimise la fréquence des points de contrôle en fonction de la taille des données et de la charge de travail. La fréquence du point de contrôle est susceptible d'être modifiée sans préavis.
Restaurer une table à un état antérieur
Utilisez la commande RESTORE pour restaurer une table à une version ou un Timestamp précédent, y compris pour les scénarios suivants :
- Vous pouvez restaurer une table déjà restaurée.
- Vous pouvez restaurer une table clonée.
Considérez les exigences suivantes :
- Pour restaurer une table, vous devez disposer de l'autorisation
MODIFYpour cette table. - Une fois les fichiers de données supprimés, manuellement ou par
VACUUM, vous ne pouvez pas restaurer une table vers une version antérieure qui référence ces fichiers. La restauration partielle à cette version est toujours possible sispark.sql.files.ignoreMissingFilesest défini surtrue. - Pour restaurer par Timestamp, utilisez les formats
yyyy-MM-dd HH:mm:ssouyyyy-MM-dd.
RESTORE TABLE target_table TO VERSION AS OF <version>;
RESTORE TABLE target_table TO TIMESTAMP AS OF <timestamp>;
Pour plus de détails sur la syntaxe, consultez RESTORE.
Comportement de streaming
La restauration est une opération de modification des données et pourrait entraîner des données dupliquées pour les charges de travail en aval. Les entrées de logs ajoutées par la commande RESTORE contiennent dataChange défini sur vrai.
Pour les charges de travail en aval, comme un job de streaming structuré qui traite les mises à jour d'une table, les entrées de journal de modification de données ajoutées par l'opération de restauration sont considérées comme de nouvelles mises à jour de données, et leur traitement peut entraîner des données dupliquées.
Par exemple :
Version de table | Opérations | Mises à jour des logs | Mises à jour des enregistrements dans le log des modifications des données |
|---|---|---|---|
0 |
|
| (name = Viktor, age = 29), (name = George, age = 55) |
1 |
|
| (nom = George, âge = 39) |
2 |
|
| Aucun enregistrement. La compaction |
3 |
|
| (nom = Viktor, âge = 29), (nom = George, âge = 55), (nom = George, âge = 39) |
Dans l'exemple précédent, la commande RESTORE entraîne des mises à jour qui avaient été vues précédemment lors de la lecture des versions 0 et 1 de la table. Si une query de streaming lit à nouveau cette table, ces fichiers sont considérés comme des données nouvellement ajoutées et sont traités à nouveau.
Restaurer les métriques
Une fois terminé, RESTORE signale les métriques suivantes sous la forme d'un DataFrame à une seule ligne :
-
table_size_after_restore: La taille de la table après restauration. -
num_of_files_after_restore: Le nombre de fichiers dans la table après la restauration. -
num_removed_files: Nombre de fichiers supprimés (logiquement) de la table. -
num_restored_files: Nombre de fichiers restaurés en raison d'un retour en arrière. -
removed_files_size: Taille totale en octets des fichiers qui sont supprimés de la table. -
restored_files_size: Taille totale en octets des fichiers restaurés.
Recherchez la dernière version de commit
Pour obtenir le numéro de version du dernier commit écrit par le SparkSession actuel sur tous les threads et toutes les tables, interrogez la configuration SQL spark.databricks.<format>.lastCommitVersionInSession. Remplacez <format> par delta ou iceberg, selon le format de votre table.
Par exemple :
- SQL
- Python
- Scala
SET spark.databricks.delta.lastCommitVersionInSession
spark.conf.get("spark.databricks.delta.lastCommitVersionInSession")
spark.conf.get("spark.databricks.delta.lastCommitVersionInSession")
Si aucun commit n'a été effectué par le SparkSession, l'interrogation de la clé renvoie une valeur vide.
Si vous partagez le même SparkSession sur plusieurs threads, c'est similaire au partage d'une variable sur plusieurs threads. Vous pourriez rencontrer des conditions de concurrence pour les mises à jour simultanées de la valeur de configuration.