Примечание.
Для доступа к этой странице требуется авторизация. Вы можете попробовать войти или изменить каталоги.
Для доступа к этой странице требуется авторизация. Вы можете попробовать изменить каталоги.
Для таблиц Apache Iceberg и Delta Lake каждая операция, которая изменяет таблицу, создает новую версию таблицы. Используйте сведения об истории для аудита операций, отката таблицы к предыдущему состоянию или выполнения запроса к таблице на определённый момент времени с помощью time travel.
Note
Не используйте историю таблиц как долгосрочное решение для резервного копирования для архивирования данных. Используйте только последние 7 дней для операций перемещения по времени, если вы не установили конфигурации хранения данных и журналов в большее значение.
Получение истории таблицы
DESCRIBE HISTORY Выполните команду, чтобы получить сведения, включая операции, пользователя и метку времени для каждой записи в таблицу. Операции возвращаются в обратном хронологическом порядке.
Для столбцов, которые DESCRIBE HISTORY возвращают, значений в operationParameters столбце и метрик для каждой операции в operationMetrics столбце см. Схема истории таблицы и метрики операций.
Хранение журнала таблиц определяется параметром logRetentionDurationтаблицы, который составляет 30 дней по умолчанию.
Note
Управление перемещением во времени и историей таблиц осуществляется различными порогами удерживания. См. раздел "Время путешествия".
DESCRIBE HISTORY table_name -- get the full history of the table
DESCRIBE HISTORY table_name LIMIT 1 -- get the last operation only
Сведения о синтаксисе Spark SQL см. в DESCRIBE HISTORY.
Сведения о синтаксисе Scala, Java и Python см. в документации по API Delta Lake.
Обозреватель каталогов визуально отображает журнал таблиц на вкладке "Журнал ".
Определение типа OPTIMIZE операции
Автоматическое сжатие, кластеризация жидкости и порядок Z отображаются в журнале таблиц в виде OPTIMIZE операций. Чтобы определить, какой из них был запущен, проверьте столбец operationParameters.
Чтобы классифицировать каждую OPTIMIZE операцию в журнале таблицы, выполните следующее:
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;
В следующих разделах подробно описано каждое operationParameters значение. Для определений operationMetrics ключей, выбранных предыдущим запросом, см. Метрики операций.
Автоматическое сжатие
Автоматическое сжатие задает для параметра auto значение true. Azure Databricks автоматически запускает автоматическое сжатие после записи. Когда значение auto — false, команду OPTIMIZE выполнил пользователь или запланированное задание.
Например, операция автоматического сжатия показывает следующее:
operationParameters: {
"auto": "true"
}
Дополнительные сведения об автоматическом уплотнении см. в разделе "Автоматическое сжатие".
Кластеризация жидкости
Кластеризация Liquid заполняет параметр clusterBy именами столбцов кластеризации.
clusterBy Пустой массив ([]) указывает только сжатие файлов.
Например, операция с кластеризованными данными по столбцам dateregion показывает следующее:
operationParameters: {
"clusterBy": "[\"date\",\"region\"]"
}
Дополнительные сведения о кластеризации жидкости см. в разделе "Использование кластеризации жидкости" для таблиц.
Z-порядок
Z-упорядочение заполняет zOrderBy параметр именами столбцов Z-порядка. Пустой zOrderBy массив ([]) указывает, что операция не применяла порядок Z.
Например, операция, применяющая порядок Z в столбце date , показывает следующее:
operationParameters: {
"zOrderBy": "[\"date\"]"
}
Область операций
Параметр predicate указывает, запущена ли операция в полной таблице или только в ней:
-
predicateПустой массив ([]) означает, что операция выполняется во всей таблице. - Заполненный
predicateмассив означает, что целеваяOPTIMIZE table_name WHERE <partition_predicate>команда выполнялась только в разделах, которые соответствуют предикату.
Например, операция, выполняемая для разделов, соответствующих year = 2024, показывает следующее:
operationParameters: {
"predicate": "[\"'year = 2024\"]"
}
Переход по времени
Путешествие во времени поддерживает выполнение запросов к предыдущим версиям таблиц на основе временной метки или версии таблицы (как это записано в журнале транзакций). Для таких приложений можно использовать поездку по времени:
- Повторное создание анализа, отчетов или выходных данных, таких как выходные данные модели машинного обучения. Это может быть полезно для отладки или аудита, особенно в регулируемых отраслях.
- Написание сложных темпоральных запросов.
- Устранение ошибок в данных.
- Обеспечение изоляции моментальных снимков для набора запросов при работе с таблицами с частыми изменениями.
Note
В Databricks Runtime 18.0 и более поздних версиях запросы с временной привязкой блокируются, если запрашивается версия, которая раньше deletedFileRetentionDuration свойства таблицы (по умолчанию 7 дней). Для управляемых таблиц каталога Unity это относится к Databricks Runtime 12.2 и выше.
Синтаксис перемещения по времени
Чтобы запросить таблицу со временем, добавьте предложение после спецификации имени таблицы.
-
timestamp_expressionможет быть одним из следующих вариантов:-
'2018-10-18T22:15:12.013Z', то есть строка, которая может преобразовываться в метку времени. cast('2018-10-18 13:36:32 CEST' as timestamp)-
'2018-10-18', то есть строка с датой. current_timestamp() - interval 12 hoursdate_sub(current_date(), 1)- Любое другое выражение, которое является меткой времени или может быть преобразовано в неё
-
-
version— это длинное значение, которое можно получить из выходных данныхDESCRIBE HISTORY table_spec.
Ни timestamp_expression, ни version не может быть подзапросом.
Принимаются только строки метки даты или времени. Например, "2019-01-01" и "2019-01-01T00:00:00.000Z". См. следующий код для примера синтаксиса:
SQL
SELECT * FROM people10m TIMESTAMP AS OF '2018-10-18T22:15:12.013Z';
SELECT * FROM people10m VERSION AS OF 123;
Python
df1 = spark.read.option("timestampAsOf", "2019-01-01").table("people10m")
df2 = spark.read.option("versionAsOf", 123).table("people10m")
Можно также использовать @ синтаксис, чтобы указать метку времени или версию в составе имени таблицы. Метка времени должна быть указана в формате yyyyMMddHHmmssSSS. Можно указать версию с @v. См. следующий код для примера синтаксиса:
SQL
-- Timestamp version
SELECT * FROM people10m@20190101000000000
-- Version number
SELECT * FROM people10m@v123
Python
# Timestamp version
spark.read.table("people10m@20190101000000000")
# Version number
spark.read.table("people10m@v123")
Настройка хранения данных для запросов на поездки по времени
Чтобы выполнить запрос к предыдущей версии таблицы, необходимо сохранить как журнал, так и файлы данных для этой версии:
- Файлы данных удаляются, когда
VACUUMзапускается для таблицы. - Файлы журнала автоматически удаляются после создания контрольных точек для версий таблицы.
Чтобы увеличить порог хранения данных для таблиц, необходимо настроить следующие свойства таблиц, заменив <format> на delta или iceberg:
-
<format>.logRetentionDuration = "interval <interval>": управляет продолжительностью хранения истории для таблицы. Значение по умолчанию —interval 30 days.- В Databricks Runtime 18.0 и более поздних версиях
logRetentionDurationдолжно быть больше или равноdeletedFileRetentionDuration. Для управляемых таблиц каталога Unity это относится к Databricks Runtime 12.2 и выше.
- В Databricks Runtime 18.0 и более поздних версиях
-
<format>.deletedFileRetentionDuration = "interval <interval>": определяет пороговое значениеVACUUM, используемое для удаления файлов данных, на которые больше не ссылается текущая версия таблицы. Значение по умолчанию —interval 7 days.
Например, чтобы получить доступ к 30 дням исторических данных, задайте delta.deletedFileRetentionDuration = "interval 30 days"значение, соответствующее параметру delta.logRetentionDurationпо умолчанию.
Important
Увеличение порогового значения хранения данных может привести к увеличению затрат на хранение, так как сохраняются дополнительные файлы данных.
Вы можете указать свойства таблицы во время создания таблицы или задать их с помощью инструкции ALTER TABLE . См. справочник по свойствам таблицы.
Примеры путешествия по времени
Чтобы восстановить случайно удалённые данные в таблице пользователя 111:
INSERT INTO my_table
SELECT * FROM my_table TIMESTAMP AS OF date_sub(current_date(), 1)
WHERE userId = 111
Чтобы исправить случайные неправильные обновления таблицы, выполните следующие действия.
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 *
Чтобы запросить количество новых клиентов, добавленных на прошлой неделе:
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
Контрольные точки журнала транзакций
Журнал транзакций записывает версии таблиц в виде JSON-файлов в каталоге журнала транзакций вместе с данными таблицы.
Для оптимизации запросов контрольных точек версии таблиц агрегируются в файлы контрольных точек Parquet, что повышает производительность, предотвращая необходимость считывания всех версий журнала таблиц JSON. Пользователям не нужно взаимодействовать с контрольными точками напрямую.
Azure Databricks оптимизирует частоту контрольных точек для размера данных и рабочей нагрузки. Частота контрольных точек подлежит изменению без уведомления.
Восстановление таблицы в более раннее состояние
RESTORE Используйте команду для восстановления таблицы до предыдущей версии или метки времени, включая следующие сценарии:
- Вы можете восстановить уже восстановленную таблицу.
- Вы можете восстановить клонированную таблицу.
Рассмотрите следующие требования:
- Чтобы восстановить таблицу, необходимо иметь
MODIFYразрешение для таблицы. - После удаления файлов данных вручную или с помощью
VACUUMневозможно восстановить таблицу до более ранней версии, которая ссылается на эти файлы. Восстановление до этой версии по-прежнему возможно, если дляspark.sql.files.ignoreMissingFilesзадано значениеtrue. - Чтобы восстановить по метке времени, используйте форматы
yyyy-MM-dd HH:mm:ssилиyyyy-MM-dd.
RESTORE TABLE target_table TO VERSION AS OF <version>;
RESTORE TABLE target_table TO TIMESTAMP AS OF <timestamp>;
Сведения о синтаксисе см. в RESTORE.
Поведение потоковой передачи
Восстановление — это операция, изменяющая данные, и она может привести к дублированию данных в последующих рабочих нагрузках. Записи журнала, добавленные командой RESTORE, содержат dataChange, установленный в значение true.
Для последующих процессов, таких как Structured Streaming-задание, которое обрабатывает обновления в таблице, записи журнала изменений данных, добавленные операцией восстановления, считаются новыми обновлениями данных, и их обработка может привести к дублированию данных.
Рассмотрим пример.
| Версия таблицы | Операция | Обновления журналов | Записи в обновлениях журнала изменений данных |
|---|---|---|---|
| 0 | INSERT |
AddFile(/path/to/file-1, dataChange = true) |
(имя = Виктор, возраст = 29), (имя = Джордж, возраст = 55) |
| 1 | INSERT |
AddFile(/path/to/file-2, dataChange = true) |
(имя = Джордж, возраст = 39) |
| 2 | OPTIMIZE |
AddFile(/path/to/file-3, dataChange = false), RemoveFile(/path/to/file-1), RemoveFile(/path/to/file-2) |
Нет записей.
OPTIMIZE Сжатие не изменяет данные в таблице. |
| 3 | RESTORE(version=1) |
RemoveFile(/path/to/file-3), AddFile(/path/to/file-1, dataChange = true), AddFile(/path/to/file-2, dataChange = true) |
(имя = Виктор, возраст = 29), (имя = Джордж, возраст = 55), (имя = Джордж, возраст = 39) |
В предыдущем примере RESTORE команда приводит к обновлениям, которые ранее были замечены при чтении таблицы версии 0 и 1. Если потоковый запрос снова считывает эту таблицу, эти файлы считаются только что добавленными данными и обрабатываются снова.
Восстановить метрики
После завершения RESTORE выводит следующие метрики в виде DataFrame из одной строки:
table_size_after_restore— размер таблицы после восстановления;num_of_files_after_restore— число файлов в таблице после восстановления;num_removed_files— число файлов, удаленных (логически удаленных) из таблицы;num_restored_files— число файлов, восстановленных из-за отката;removed_files_size— общий размер в байтах файлов, удаленных из таблицы;restored_files_size— общий размер восстановленных файлов в байтах.
Найти последнюю версию коммита
Чтобы получить номер версии последней фиксации, записанной текущим SparkSession, для всех потоков и всех таблиц, выполните запрос конфигурации SQL spark.databricks.<format>.lastCommitVersionInSession. Замените <format> либо deltaicebergв зависимости от формата таблицы.
Рассмотрим пример.
SQL
SET spark.databricks.delta.lastCommitVersionInSession
Python
spark.conf.get("spark.databricks.delta.lastCommitVersionInSession")
Scala
spark.conf.get("spark.databricks.delta.lastCommitVersionInSession")
Если SparkSession не сделал коммиты, запрос ключа возвращает пустое значение.
Note
Если один и тот же SparkSession используется несколькими потоками, это похоже на совместное использование переменной несколькими потоками. Вы можете столкнуться с условиями гонки для параллельных обновлений значения конфигурации.