Использовать канал данных изменений в Azure Databricks

Канал изменений данных (CDF) отслеживает изменения на уровне строк между версиями таблицы Delta Lake или таблицы Apache Iceberg v3. Каждая запись изменения включает данные строк вместе с метаданными, указывающими, была ли строка вставлена, обновлена или удалена.

Вы можете использовать веб-канал изменений для распространенных вариантов использования данных, в том числе:

  • Добавочные конвейеры ETL, обрабатывающие только строки, которые изменились с момента последнего запуска конвейера.
  • Журналы аудита, отслеживающие изменения данных для соблюдения требований соответствия и управления данными.
  • Задачи репликации данных, которые синхронизируют изменения с целевыми таблицами, кэшами или внешними системами.

Azure Databricks поддерживает два подхода:

  • Автоматическая лента изменений данных: Определяет изменения при чтении таблицы с помощью метаданных происхождения строк. Для этого не требуется отдельная конфигурация таблицы и работает в таблицах Delta Lake и Apache Iceberg версии 3. См. автоматическую передачу данных об изменениях.
  • Устаревший поток данных об изменениях: материализует изменения при записи в таблицу. Поддерживает только таблицы Delta Lake. Требуется отдельная конфигурация таблицы. См. устаревший поток данных об изменениях для Delta Lake.

Канал данных об автоматических изменениях

Автоматический поток данных изменений вычисляет изменения на уровне строк во время запроса, а не в момент записи, используя отслеживание строк в таблицах Delta Lake и линии строк в таблицах Apache Iceberg v3. В отличие от устаревшего механизма Change Data Feed, не требуется индивидуальная настройка таблиц. Любая таблица, соответствующая требованиям, поддерживает его автоматически. См. раздел Требования.

Поскольку изменения не вычисляются при каждой записи для операций MERGE INTO и UPDATE, автоматический канал данных об изменениях повышает производительность записи и снижает затраты на хранение по сравнению с прежним каналом данных об изменениях.

Автоматический канал данных об изменениях использует те же API table_changes() и readChangeFeed, что и устаревший канал данных об изменениях, и работает с пакетными запросами, Structured Streaming и Databricks-to-Databricks OpenSharing. См. Изменения в пакетных запросах и Инкрементальная обработка данных об изменениях.

Требования

  • Databricks Runtime 19 и выше
  • Поддерживаемый формат таблицы, зарегистрированный в каталоге Unity:
    • Управляемая таблица в формате Delta Lake с включенным отслеживанием строк или в формате Iceberg версии 3.
    • Внешняя таблица в формате Delta Lake с включенным отслеживанием строк.

См. типы таблиц каталога Databricks Unity.

Замечание

Поток данных об изменениях не является частью спецификации Apache Iceberg. Средства чтения Azure Databricks могут запрашивать автоматически формируемый поток данных об изменениях для таблиц Apache Iceberg v3, но внешние средства чтения Iceberg не могут. См. спецификацию таблицы Айсберга.

Для Delta Lake только средства чтения Azure Databricks могут выполнять запросы к автоматическому потоку данных об изменениях.

Использование канала изменений данных

Чтобы использовать канал передачи данных об изменениях, убедитесь, что вы используете таблицу, соответствующую требованиям. См. раздел Требования.

Чтобы выполнить пакетное чтение потока изменений данных, выполните следующие действия.

Python

spark.read \
  .option("readChangeFeed", "true") \
  .option("startingVersion", 0) \
  .table("<table_name>")

Scala

spark.read
  .option("readChangeFeed", "true")
  .option("startingVersion", 0)
  .table("<table_name>")

SQL

SELECT * FROM table_changes('<table_name>', 0)

Дополнительные сведения о пакетном чтении для ленты данных об изменениях см. в разделе "Чтение изменений в пакетных запросах".

Чтобы считывать поток изменённых данных в потоковом режиме, выполните следующие действия:

Python

(spark.readStream
  .option("readChangeFeed", "true")
  .table("<table_name>")
)

Scala

spark.readStream
  .option("readChangeFeed", "true")
  .table("<table_name>")

Дополнительные сведения о чтении в потоковом режиме для потока данных об изменениях см. в разделе "Инкрементная обработка данных об изменениях".

Переход с устаревшего потока данных об изменениях

Чтобы перевести таблицу Delta Lake из устаревшего потока данных об изменениях в автоматический поток данных об изменениях, выполните следующие действия:

  1. Убедитесь, что таблица соответствует требованиям.
  2. Отключите устаревший поток изменённых данных, выполнив следующую команду:
ALTER TABLE <table_name> UNSET TBLPROPERTIES ('delta.enableChangeDataFeed');

Нельзя одновременно использовать и устаревшие, и автоматические потоки данных об изменениях.

Изменение схемы канала данных

При чтении из канала данных об изменениях для таблицы запрос использует схему последней версии таблицы. Azure Databricks поддерживает большинство операций изменения схемы и эволюции, но таблицы с сопоставлением столбцов имеют ограничения. См. таблицы с сопоставлением столбцов.

Помимо столбцов данных из схемы таблицы Delta Lake, канал изменений содержит столбцы метаданных, определяющие тип события изменения:

Название столбца Type Ценности
_change_type String Содержит: insert, , update_preimageupdate_postimagedelete.
preimage значением перед обновлением postimage является значение после обновления.
_commit_version Long Содержит: разностный журнал или версия таблицы, содержащая изменение.
_commit_timestamp Timestamp Содержит: метку времени, связанную с моментом создания коммита.

Если схема содержит столбцы с теми же именами, что и у этих столбцов метаданных, нельзя использовать канал отслеживания изменений на таблице. Прежде чем включить канал передачи изменённых данных, переименуйте столбцы в таблице, чтобы устранить этот конфликт.

Добавочная обработка данных об изменении

Databricks рекомендует использовать канал отслеживания изменений данных в сочетании с Structured Streaming для инкрементальной обработки изменений в таблицах. Чтобы автоматически отслеживать версии канала изменений данных вашей таблицы, необходимо использовать структурированную потоковую обработку в Azure Databricks. Сведения об обработке CDC с таблицами типа 1 или типа 2 см. в API AUTO CDC: упрощение захвата измененных данных с помощью конвейеров.

Когда поток впервые запускается, поток передачи данных изменений возвращает актуальный снимок таблицы как записи INSERT, а затем возвращает последующие изменения как данные об изменениях. Потоки данных об изменениях одновременно записывают в журнал транзакций таблицы как данные об изменениях, так и новые строки данных.

Чтобы настроить поток для чтения канала данных об изменениях таблицы, задайте для параметра readChangeFeed значение true следующим образом:

Python

(spark.readStream
  .option("readChangeFeed", "true")
  .table("myTable")
)

Scala

spark.readStream
  .option("readChangeFeed", "true")
  .table("myTable")

Ограничения скорости

Azure Databricks поддерживает ограничения скорости (maxFilesPerTrigger, maxBytesPerTriggerи excludeRegex при чтении измененных данных). Чтобы просмотреть полный список параметров потоковой передачи Delta Lake, см. раздел Delta Lake.

При необходимости можно указать начальную версию, см. раздел "Указание начальной версии". Для версий, отличных от исходного снимка, ограничения по скорости применяются атомарно к коммитам целиком. Либо текущий пакет содержит весь коммит, либо текущий пакет переносит коммит в следующий пакет.

Журнал таблиц воспроизведения

Канал передачи данных об изменениях не предназначен для того, чтобы служить постоянной записью всех изменений в таблице. Записываются только изменения, произошедшие после включения канала данных об изменениях. Вы можете начать новое потоковое чтение, чтобы получить текущую версию и все последующие изменения.

Записи в потоке измененных данных являются временными и доступны только в течение указанного периода хранения. Журналы транзакций удаляют версии таблиц и соответствующие им версии ленты изменений данных через регулярные промежутки времени. При удалении версии вы больше не сможете прочитать веб-канал измененных данных для этой версии.

Архивирование данных об изменениях для постоянного хранения в журнале

Если в вашем случае использования требуется сохранить постоянную историю всех изменений в таблице, используйте добавочную логику для записи записей из канала изменений в новую таблицу.

В следующем примере показано использование trigger.AvailableNow для обработки доступных данных в виде пакетной рабочей нагрузки для аудита или полного воспроизведения изменений:

Python
(spark.readStream
  .option("readChangeFeed", "true")
  .table("source_table")
  .writeStream
  .option("checkpointLocation", <checkpoint-path>)
  .trigger(availableNow=True)
  .toTable("target_table")
)
Scala
spark.readStream
  .option("readChangeFeed", "true")
  .table("source_table")
  .writeStream
  .option("checkpointLocation", <checkpoint-path>)
  .trigger(Trigger.AvailableNow)
  .toTable("target_table")

Указание начальной версии

Чтобы считывать изменения из определенной точки, укажите начальную версию с помощью метки времени или номера версии. Начальные версии данных необходимы для пакетных операций чтения. При необходимости можно указать конечную версию, чтобы ограничить диапазон. Дополнительные сведения об истории таблицы см. в разделе "Путешествие во времени".

При настройке рабочих нагрузок Structured Streaming, использующих канал передачи данных об изменениях, указание начальной версии может повлиять на производительность обработки:

  • Новые конвейеры обработки данных обычно используют поведение по умолчанию, которое записывает все существующие записи в таблице как INSERT операции при первом запуске потока.
  • Если целевая таблица уже содержит все записи с соответствующими изменениями до определенной точки, укажите начальную версию, чтобы избежать обработки состояния исходной таблицы в виде INSERT событий.

В следующем примере показано, как восстановиться после сбоя потоковой передачи с поврежденной контрольной точкой. В этом примере предполагается следующее:

  1. Канал изменений данных включен в исходной таблице при создании таблицы.
  2. Целевая нижестоящая таблица обработала все изменения до версии 75 включительно.
  3. Журнал версий исходной таблицы доступен для версий 70 и выше.

При определении потока записи в существующую целевую таблицу необходимо указать новое расположение контрольной точки:

Python

(spark.readStream
  .option("readChangeFeed", "true")
  .option("startingVersion", 76)
  .table("source_table")
  .writeStream
  .option("checkpointLocation", "<new-checkpoint-path>")
  .toTable("target_table")
)

Scala

spark.readStream
  .option("readChangeFeed", "true")
  .option("startingVersion", 76)
  .table("source_table")
  .writeStream
  .option("checkpointLocation", "<new-checkpoint-path>")
  .toTable("target_table")

Important

Если в журнале таблиц указана начальная версия и эта версия недоступна, поток не может начинаться с новой контрольной точки. Так как управляемые таблицы автоматически очищают исторические версии, все указанные начальные версии в конечном итоге удаляются.

См. историю таблицы реплея.

Чтение изменений в пакетных запросах

С помощью синтаксиса пакетного запроса можно считывать все изменения, начиная с конкретной версии или считывать изменения в указанном диапазоне версий следующим образом:

  • Укажите версии в виде целых чисел и меток времени в виде строк в формате yyyy-MM-dd[ HH:mm:ss[.SSS]].
  • Начальные и конечные версии являются включающими. Чтобы прочитать исходную версию до последней версии, укажите только начальную версию.
  • Если указать версию, предшествующую включению потока изменённых данных, возникнет ошибка.

Чтобы использовать пакетные операции чтения с параметрами начальной и конечной версии, сделайте следующее:

SQL

Чтобы прочитать информацию с версии 0 по 10, выполните следующие действия:

SELECT * FROM table_changes('tableName', 0, 10)

Чтобы прочитать между двумя версиями метки времени, сделайте следующее:

--
SELECT * FROM table_changes('tableName', '2021-04-21 05:45:46', '2021-05-21 12:00:00')

Чтобы прочитать исходную версию до последней, сделайте следующее:

SELECT * FROM table_changes('tableName', 0)

Чтобы прочитать изменения для таблицы со специальными символами в имени, сделайте следующее:

SELECT * FROM table_changes('`schema`.`dotted.tableName`', '2021-04-21 06:45:46', '2021-05-21 12:00:00')

См. табличную функцию table_changes.

Python

Чтобы прочитать информацию с версии 0 по 10, выполните следующие действия:

spark.read \
  .option("readChangeFeed", "true") \
  .option("startingVersion", 0) \
  .option("endingVersion", 10) \
  .table("myDeltaTable")

Чтобы прочитать данные между двумя метками времени, выполните следующие действия:

spark.read \
  .option("readChangeFeed", "true") \
  .option("startingTimestamp", '2021-04-21 05:45:46') \
  .option("endingTimestamp", '2021-05-21 12:00:00') \
  .table("myDeltaTable")

Чтобы прочитать исходную версию до последней, сделайте следующее:

spark.read \
  .option("readChangeFeed", "true") \
  .option("startingVersion", 0) \
  .table("myDeltaTable")

Scala

Чтобы прочитать информацию с версии 0 по 10, выполните следующие действия:

spark.read
  .option("readChangeFeed", "true")
  .option("startingVersion", 0)
  .option("endingVersion", 10)
  .table("myDeltaTable")

Чтобы прочитать данные между двумя метками времени, выполните следующие действия:

spark.read
  .option("readChangeFeed", "true")
  .option("startingTimestamp", "2021-04-21 05:45:46")
  .option("endingTimestamp", "2021-05-21 12:00:00")
  .table("myDeltaTable")

Чтобы прочитать исходную версию до последней, сделайте следующее:

spark.read
  .option("readChangeFeed", "true")
  .option("startingVersion", 0)
  .table("myDeltaTable")

Обработка версий вне диапазона

По умолчанию, если указать версию или временную метку, превышающую последний коммит, запрос вернет ошибку timestampGreaterThanLatestCommit.

В Databricks Runtime 11.3 LTS и более поздних версиях можно включить терпимость к версиям вне диапазона следующим образом:

SET spark.databricks.delta.changeDataFeed.timestampOutOfRange.enabled = true;

Если эта конфигурация включена, запрос возвращает различные результаты следующим образом:

  • Начальная версия или метка времени позже последнего коммита возвращает пустой результат.
  • Конечная версия или метка времени за пределами последней фиксации возвращает все изменения от начала до последней фиксации.

Прежний поток данных об изменениях для Delta Lake

Устаревший поток данных об изменениях требует ручной настройки отдельных таблиц Delta Lake. Поскольку поток данных об изменениях не включен в спецификацию Apache Iceberg, таблицы Apache Iceberg не поддерживаются. Databricks рекомендует перейти на автоматический канал передачи данных об изменениях. См. статью "Миграция из устаревшего канала данных об изменениях".

Когда включён устаревший канал передачи данных об изменениях, среда выполнения записывает события изменения для всех данных, записываемых в таблицу. К ним относятся данные строк вместе с метаданными, которые указывают, была ли указана строка вставлена, удалена или обновлена.

Устаревший поток данных об изменениях использует те же API чтения readChangeFeed и table_changes(), что и автоматический поток данных об изменениях. См. Инкрементная обработка данных об изменениях и Чтение изменений в пакетных запросах.

Включите устаревший канал данных об изменениях

Необходимо явно включить прежний поток изменений данных для отдельных таблиц. Используйте один из следующих методов:

Новая таблица

Задайте свойство delta.enableChangeDataFeed = true таблицы в команде CREATE TABLE .

CREATE TABLE student (id INT, name STRING, age INT)
  TBLPROPERTIES (delta.enableChangeDataFeed = true)

Замечание

Если вы отключите прежний канал передачи данных об изменениях на какое-то время, а затем снова включите его, данные за этот промежуток будет невозможно запросить. Используйте автоматический поток данных об изменениях для получения сведений об изменениях за указанный интервал. См. автоматическую передачу данных об изменениях.

Существующая таблица

Задайте свойство delta.enableChangeDataFeed = true таблицы в команде ALTER TABLE .

ALTER TABLE myDeltaTable
  SET TBLPROPERTIES (delta.enableChangeDataFeed = true)

Рекомендации по хранению

Управляемые таблицы эффективно записывают изменения данных и могут использовать другие функции для оптимизации макета хранилища.

При использовании устаревшего канала данных об изменении необходимо учитывать следующее поведение хранилища:

  • Возможно, вы увидите небольшое увеличение затрат на хранение, так как изменения могут быть записаны в отдельных файлах.
  • Некоторые операции, такие как операции только вставки или удаление целых разделов, не создают файлы данных об изменениях. Azure Databricks вычисляет поток данных об изменениях непосредственно на основе журнала транзакций.
  • Файлы данных об изменениях используют политику хранения таблицы. Команда VACUUM удаляет файлы измененных данных, а изменения из журнала транзакций используют политику хранения контрольных точек.

Databricks рекомендует не пытаться воссоздать поток данных об изменениях, напрямую запрашивая файлы данных об изменениях. Всегда используйте API Delta Lake и Apache Iceberg.

Ограничения

Учитывайте следующие ограничения для потоков данных об изменениях:

Таблицы с сопоставлением столбцов

С включенным сопоставлением столбцов в таблице Delta Lake можно удалить или переименовать столбцы без перезаписи файлов данных. См. как переименовывать и удалять столбцы с использованием сопоставления столбцов в Delta Lake.

Однако потоки данных об изменениях имеют ограничения после неаддитивных изменений схемы. Изменения схемы, не являющиеся аддитивными, включают следующие операции:

  • Переименование или удаление столбцов.
  • Изменение типов данных столбца.
  • Изменение допустимости значения NULL для столбца, например с помощью ALTER COLUMN ... SET NOT NULL. См. NOT NULLограничение.

Невозможно считывать каналы данных об изменениях для транзакции или диапазона, где имеет место неаддитивное изменение схемы.

Чтобы разрешить неаддитивные изменения схемы до или после указанного диапазона пакетных операций чтения, запросы используют схему конечной версии диапазона, а не последнюю версию таблицы. Запросы по-прежнему завершаются ошибкой, если диапазон версий охватывает неаддитивное изменение схемы.

Автоматический поток данных об изменениях

  • Поскольку автоматический канал данных об изменениях не поддерживается в спецификации Apache Iceberg, внешние клиенты Iceberg не могут выполнять к нему запросы. См. спецификацию таблицы Айсберга.
  • Для транзакций, состоящих из нескольких инструкций, автоматический канал передачи данных об изменениях не поддерживается, если исходная таблица была изменена во время транзакции.
  • Автоматическая лента передачи данных об изменениях не поддерживается для таблиц с фильтрами по строкам или масками столбцов. См. фильтры строк и маски столбцов.
  • Запросы к каналу изменений данных не могут охватывать версии таблицы, в которых произошло неаддитивное изменение схемы, например переименование столбца, его удаление или изменение типа данных. Разделить запрос на диапазоны до и после изменения схемы.