Примечание.
Для доступа к этой странице требуется авторизация. Вы можете попробовать войти или изменить каталоги.
Для доступа к этой странице требуется авторизация. Вы можете попробовать изменить каталоги.
Конвейеры Lakeflow упрощают захват изменений данных (CDC) с помощью API AUTO CDC и AUTO CDC FROM SNAPSHOT. Эти API автоматизируют сложность вычислений медленно изменяющихся измерений (SCD) типа 1 и типа 2 из потока CDC (отслеживания изменений данных) или моментальных снимков базы данных.
AUTO CDC API также поддерживает двухвременное отслеживание, которое фиксирует изменения в двух временных измерениях (бета-версия). Дополнительные сведения о SCD Type 1 и Type 2 см. в разделе Захват изменений данных и моментальные снимки. Дополнительные сведения о отслеживании укусов см. в разделе Bitmporal AUTO CDC.
Замечание
AUTO CDC API-интерфейсы заменяют APPLY CHANGES API и имеют тот же синтаксис.
APPLY CHANGES API по-прежнему доступны, но Databricks рекомендует использовать AUTO CDC API на их месте.
Используемый API зависит от источника измененных данных:
-
AUTO CDC: используйте данный элемент, если в исходной базе данных включена поддержка CDC.AUTO CDCобрабатывает изменения из канала данных изменений (CDF). Поддерживается как в конвейерных интерфейсах SQL, так и в Python. -
AUTO CDC FROM SNAPSHOT: используйте этот параметр, если CDC не включен в исходной базе данных и доступны только моментальные снимки. Этот API сначала сравнивает моментальные снимки, чтобы определить изменения, а затем обрабатывает их. Она поддерживается только в интерфейсе Python.
Оба API поддерживают обновление таблиц с помощью SCD Type 1 и Type 2:
- Используйте SCD Type 1 для непосредственного обновления записей. История обновленных записей не сохраняется.
- Используйте SCD Type 2 для сохранения истории записей при всех обновлениях или обновлении указанного набора столбцов.
Только для AUTO CDC можно также использовать битемпоральное хранилище, которое расширяет историю SCD типа 2 и позволяет отслеживать изменения по двум временным измерениям: бизнес-время и системное время. Bitmporal находится в бета-версии. См. статью Bitmporal AUTO CDC.
AUTO CDC также поддерживает частичные обновления, в которых запись изменений обновляет только подмножество столбцов. См. раздел "Применить частичные обновления".
AUTO CDC API не поддерживаются декларативными конвейерами Apache Spark.
Сведения о синтаксисе и других ссылках см. в AUTO CDC INTO (конвейеры), create_auto_cdc_flow и create_auto_cdc_from_snapshot_flow.
Замечание
На этой странице описывается обновление таблиц в конвейерах на основе изменений в исходных данных. Сведения о записи и запросе сведений об изменении на уровне строк для таблиц Delta см. в разделе "Использование канала изменений" в Azure Databricks.
Требования
Чтобы использовать API CDC, ваш конвейер должен быть настроен на использование бессерверных конвейеров Lakeflow или конвейеров Lakeflow Pro или Advancedредакций.
Как работает AUTO CDC
Чтобы выполнить обработку CDC с AUTO CDC, создайте потоковую таблицу, а затем используйте оператор AUTO CDC ... INTO в SQL или функцию create_auto_cdc_flow() в Python, чтобы указать источник, ключи и последовательности для канала изменений. Сведения о том, как работает логика упорядочивания и SCD, см. в разделе "Изменение данных и моментальных снимков". Ознакомьтесь с примерами AUTO CDC.
Для первоначальной гидратации из источника с каналом изменений используйте AUTO CDC с потоком once, а затем продолжайте обработку потока изменений. См. статью "Репликация внешней таблицы RDBMS" с помощью AUTO CDC.
Сведения о синтаксисе см. в разделе AUTO CDC INTO (конвейеры) или create_auto_cdc_flow.
Как работает AUTO CDC из SNAPSHOT
AUTO CDC FROM SNAPSHOT определяет изменения исходных данных путем сравнения последовательных моментальных снимков. Он поддерживается только в интерфейсе конвейера Python. Моментальные снимки можно считывать непосредственно из таблицы Delta, файлов облачного хранилища или JDBC.
Чтобы выполнить обработку данных CDC с помощью AUTO CDC FROM SNAPSHOT, создайте потоковую таблицу, после этого используйте функцию create_auto_cdc_from_snapshot_flow() для указания моментального снимка, ключей и других аргументов. Дополнительные сведения о двух шаблонах приема данных и их использовании см. в разделе "Шаблоны обработки моментальных снимков". Ознакомьтесь с примерами AUTO CDC FROM SNAPSHOT.
Сведения о синтаксисе см. в разделе create_auto_cdc_from_snapshot_flow.
Используйте несколько столбцов для упорядочивания
Чтобы упорядочить по нескольким столбцам (например, метке времени и идентификатору для разрыва связей), используйте STRUCT, чтобы объединить их. API сортирует сначала по первому полю, а в случае равенства учитывает второе поле и так далее.
SQL
SEQUENCE BY STRUCT(timestamp_col, id_col)
Питон
sequence_by = struct("timestamp_col", "id_col")
Примеры AUTO CDC
В следующих примерах показана обработка данных SCD типов 1 и 2 с использованием источника данных о изменениях. Пример данных создает новые записи пользователей, удаляет запись пользователя и обновляет записи пользователей. В примере SCD Type 1 последние UPDATE операции приходят поздно и удаляются из целевой таблицы, иллюстрируя обработку событий в неправильном порядке.
Ниже приведены входные записи, используемые в этих примерах. Эти данные создаются путем выполнения запроса в разделе "Создание примера данных ".
| userId | имя | city | Операция | номер последовательности |
|---|---|---|---|---|
| 124 | Рауль | Оахака | INSERT | 1 |
| 123 | Isabel | Монтеррей | INSERT | 1 |
| 125 | Мерседес | Тихуана | INSERT | 2 |
| 126 | Лилия | Канкун | INSERT | 2 |
| 123 | ноль | ноль | Удалить | 6 |
| 125 | Мерседес | Guadalajara | UPDATE | 6 |
| 125 | Мерседес | Мехикали | UPDATE | 5 |
| 123 | Isabel | Чиуауа | UPDATE | 5 |
Если вы раскомментируете окончательную строку в примере запроса создания данных, она добавляет следующую запись, которая указывает на необходимость очистки таблицы в sequenceNum=3.
| userId | имя | city | Операция | номер последовательности |
|---|---|---|---|---|
| ноль | ноль | ноль | УКОРАЧИВАТЬ | 3 |
Замечание
Все приведенные ниже примеры включают параметры для указания обоих DELETE и TRUNCATE операций, но каждый из них является необязательным.
Создание примеров данных
Выполните следующие инструкции, чтобы создать образец набора данных. Этот код не предназначен для запуска в рамках определения конвейера. Запустите его из папки исследования конвейера, а не из папки преобразований.
CREATE SCHEMA IF NOT EXISTS main.cdc_tutorial;
CREATE TABLE main.cdc_tutorial.users_cdf
AS SELECT
col1 AS userId,
col2 AS name,
col3 AS city,
col4 AS operation,
col5 AS sequenceNum
FROM (
VALUES
-- Initial load.
(124, "Raul", "Oaxaca", "INSERT", 1),
(123, "Isabel", "Monterrey", "INSERT", 1),
-- New users.
(125, "Mercedes", "Tijuana", "INSERT", 2),
(126, "Lily", "Cancun", "INSERT", 2),
-- Isabel is removed from the system and Mercedes moved to Guadalajara.
(123, null, null, "DELETE", 6),
(125, "Mercedes", "Guadalajara", "UPDATE", 6),
-- This batch of updates arrived out of order. The batch at sequenceNum 6 is the final state.
(125, "Mercedes", "Mexicali", "UPDATE", 5),
(123, "Isabel", "Chihuahua", "UPDATE", 5)
-- Uncomment to test TRUNCATE.
-- ,(null, null, null, "TRUNCATE", 3)
);
Обработка обновлений SCD Type 1
SCD Type 1 сохраняет только последнюю версию каждой записи. В следующем примере считывается из потока данных изменений, созданного выше, и применяются изменения к целевой потоковой таблице. Что такое конвейеры? для выполнения этого кода.
Питон
from pyspark import pipelines as dp
from pyspark.sql.functions import col, expr
@dp.view
def users():
return spark.readStream.table("main.cdc_tutorial.users_cdf")
dp.create_streaming_table("users_current")
dp.create_auto_cdc_flow(
target = "users_current",
source = "users",
keys = ["userId"],
sequence_by = col("sequenceNum"),
apply_as_deletes = expr("operation = 'DELETE'"),
apply_as_truncates = expr("operation = 'TRUNCATE'"),
except_column_list = ["operation", "sequenceNum"],
stored_as_scd_type = 1
)
SQL
CREATE OR REFRESH STREAMING TABLE users_current;
CREATE FLOW apply_cdc AS AUTO CDC INTO
users_current
FROM
stream(main.cdc_tutorial.users_cdf)
KEYS
(userId)
APPLY AS DELETE WHEN
operation = "DELETE"
APPLY AS TRUNCATE WHEN
operation = "TRUNCATE"
SEQUENCE BY
sequenceNum
COLUMNS * EXCEPT
(operation, sequenceNum)
STORED AS
SCD TYPE 1;
После выполнения примера SCD Type 1 целевая таблица содержит следующие записи:
| userId | имя | city |
|---|---|---|
| 124 | Рауль | Оахака |
| 125 | Мерседес | Guadalajara |
| 126 | Лилия | Канкун |
Пользователь 123 (Isabel) был удален и не отображается. Пользователь 125 (Mercedes) показывает только последний город (Guadalajara), так как SCD Type 1 перезаписывает предыдущие значения. Ранее UPDATE на sequenceNum=5 было отменено, потому что поступило позднее обновление на sequenceNum=6.
После выполнения примера с раскомментированной записью TRUNCATE, таблица очищается в sequenceNum=3. Это означает, что записи 124 и 126 не находятся в таблице, а окончательная целевая таблица содержит только следующую запись:
| userId | имя | city |
|---|---|---|
| 125 | Мерседес | Guadalajara |
Обработка обновлений SCD Type 2
SCD Type 2 сохраняет полный журнал изменений, создавая новые строки для каждой версии записи, а __START_AT__END_AT также столбцы, указывающие, когда каждая версия была активной.
Питон
from pyspark import pipelines as dp
from pyspark.sql.functions import col, expr
@dp.view
def users():
return spark.readStream.table("main.cdc_tutorial.users_cdf")
dp.create_streaming_table("users_history")
dp.create_auto_cdc_flow(
target = "users_history",
source = "users",
keys = ["userId"],
sequence_by = col("sequenceNum"),
apply_as_deletes = expr("operation = 'DELETE'"),
except_column_list = ["operation", "sequenceNum"],
stored_as_scd_type = "2"
)
SQL
CREATE OR REFRESH STREAMING TABLE users_history;
CREATE FLOW apply_cdc AS AUTO CDC INTO
users_history
FROM
stream(main.cdc_tutorial.users_cdf)
KEYS
(userId)
APPLY AS DELETE WHEN
operation = "DELETE"
SEQUENCE BY
sequenceNum
COLUMNS * EXCEPT
(operation, sequenceNum)
STORED AS
SCD TYPE 2;
После запуска примера SCD Type 2 целевая таблица содержит следующие записи:
| userId | имя | city | __НАЧАТЬ_С | __END_AT |
|---|---|---|---|---|
| 123 | Isabel | Монтеррей | 1 | 5 |
| 123 | Isabel | Чиуауа | 5 | 6 |
| 124 | Рауль | Оахака | 1 | ноль |
| 125 | Мерседес | Тихуана | 2 | 5 |
| 125 | Мерседес | Мехикали | 5 | 6 |
| 125 | Мерседес | Guadalajara | 6 | ноль |
| 126 | Лилия | Канкун | 2 | ноль |
Таблица сохраняет полную историю. Пользователь 123 имеет две версии (заканчиваются последовательностью 6 при удалении). Пользователь 125 имеет три версии, показывающие изменения города. Записи с __END_AT = null в настоящее время активны.
Отслеживание подмножества столбцов с помощью SCD Type 2
По умолчанию SCD Type 2 создает новую версию при каждом изменении значения столбца. Можно указать подмножество столбцов для отслеживания, чтобы изменения других столбцов обновляли текущую версию вместо создания новой записи журнала.
Пример, который следует, исключает city столбец из отслеживания истории.
Питон
from pyspark import pipelines as dp
from pyspark.sql.functions import col, expr
@dp.view
def users():
return spark.readStream.table("main.cdc_tutorial.users_cdf")
dp.create_streaming_table("users_history")
dp.create_auto_cdc_flow(
target = "users_history",
source = "users",
keys = ["userId"],
sequence_by = col("sequenceNum"),
apply_as_deletes = expr("operation = 'DELETE'"),
except_column_list = ["operation", "sequenceNum"],
stored_as_scd_type = "2",
track_history_except_column_list = ["city"]
)
SQL
CREATE OR REFRESH STREAMING TABLE users_history;
CREATE FLOW apply_cdc AS AUTO CDC INTO
users_history
FROM
stream(main.cdc_tutorial.users_cdf)
KEYS
(userId)
APPLY AS DELETE WHEN
operation = "DELETE"
SEQUENCE BY
sequenceNum
COLUMNS * EXCEPT
(operation, sequenceNum)
STORED AS
SCD TYPE 2
TRACK HISTORY ON * EXCEPT
(city)
Так как city изменения не отслеживаются, обновления города перезаписывают текущую строку вместо создания новой версии. Целевая таблица содержит следующие записи:
| userId | имя | city | __НАЧАТЬ_С | __END_AT |
|---|---|---|---|---|
| 123 | Isabel | Чиуауа | 1 | 6 |
| 124 | Рауль | Оахака | 1 | ноль |
| 125 | Мерседес | Guadalajara | 2 | ноль |
| 126 | Лилия | Канкун | 2 | ноль |
Примеры AUTO CDC FROM SNAPSHOT
В следующих разделах приведены примеры использования AUTO CDC FROM SNAPSHOT для обработки моментальных снимков в целевых таблицах SCD Type 1 или Type 2. Информацию о том, когда использовать этот API, см. в разделе Отслеживание изменений данных и снимки данных.
Пример. Обработка моментальных снимков с помощью времени приема конвейера
Используйте этот подход, когда моментальные снимки поступают регулярно и в правильном порядке, и вы можете полагаться на метку времени выполнения конвейера для управления версиями. При каждом обновлении конвейера осуществляется прием нового моментального снимка.
Моментальные снимки можно считывать из нескольких исходных типов, включая таблицы Delta, файлы облачного хранилища и подключения JDBC.
Шаг 1. Создание примеров данных
Создайте таблицу, содержащую данные моментального снимка. Запустите следующий код из записной книжки или Databricks SQL в папке explorations конвейера:
CREATE SCHEMA IF NOT EXISTS main.cdc_tutorial;
CREATE TABLE main.cdc_tutorial.snapshot (
userId INT,
city STRING
);
INSERT INTO main.cdc_tutorial.snapshot VALUES
(1, 'Oaxaca'),
(2, 'Monterrey'),
(3, 'Tijuana');
Шаг 2. Запуск AUTO CDC FROM SNAPSHOT
Что такое конвейеры? для запуска кода на этом шаге.
Выберите тип источника для представления снимков (код создания примера генерирует таблицу Delta):
Вариант A. Чтение из таблицы Delta
from pyspark import pipelines as dp
@dp.view(name="source")
def source():
return spark.read.table("main.cdc_tutorial.snapshot")
Вариант B. Чтение из облачного хранилища
from pyspark import pipelines as dp
@dp.view(name="source")
def source():
return spark.read.format("csv").option("header", True).load("<snapshot-path>")
Вариант C. Чтение из JDBC (только для классических вычислений)
from pyspark import pipelines as dp
@dp.view(name="source")
def source():
return (spark.read
.format("jdbc")
.option("url", "<jdbc-url>")
.option("dbtable", "<table-name>")
.option("user", "<username>")
.option("password", "<password>")
.load()
)
Все параметры, запись в целевой объект
Затем добавьте целевую таблицу и поток:
dp.create_streaming_table("target")
dp.create_auto_cdc_from_snapshot_flow(
target = "target",
source = "source",
keys = ["userId"],
stored_as_scd_type = 2
)
После первого запуска конвейера все записи вставляются в виде активных строк:
| userId | city | __НАЧАТЬ_С | __END_AT |
|---|---|---|---|
| 1 | Оахака | 0 | ноль |
| 2 | Монтеррей | 0 | ноль |
| 3 | Тихуана | 0 | ноль |
Замечание
Чтобы использовать SCD Type 1 вместо Type 2 и сохранить только текущее состояние, установите stored_as_scd_type=1. В этом случае целевая таблица не включает __START_AT и __END_AT столбцы.
Шаг 3: Смоделируйте новый моментальный снимок и выполните повторный запуск
Обновите исходную таблицу, чтобы имитировать новый снимок (запустите код из ноутбука или SQL файла в папке explorations pipeline):
TRUNCATE TABLE main.cdc_tutorial.snapshot;
INSERT INTO main.cdc_tutorial.snapshot VALUES
(2, 'Carmel'),
(3, 'Los Angeles'),
(4, 'Death Valley'),
(6, 'Kings Canyon');
Повторный запуск конвейера.
AUTO CDC FROM SNAPSHOT сравнивает новый моментальный снимок с предыдущим и обнаруживает, что пользователь 1 был удален, пользователи 2 и 3 были обновлены, а пользователи 4 и 6 были вставлены. Это генерирует канал изменений и использует AUTO CDC для создания выходной таблицы.
После второго запуска с SCD Type 2 целевая таблица содержит следующие записи:
| userId | city | __НАЧАТЬ_С | __END_AT |
|---|---|---|---|
| 1 | Оахака | 0 | 1 |
| 2 | Монтеррей | 0 | 1 |
| 2 | Кармель | 1 | ноль |
| 3 | Тихуана | 0 | 1 |
| 3 | Лос-Анджелес | 1 | ноль |
| 4 | Долина смерти | 1 | ноль |
| 6 | Каньон Кингс | 1 | ноль |
Пользователь 1 был завершен (удален). Пользователи 2 и 3 имеют две версии, показывающие изменения города. Пользователи 4 и 6 были недавно добавлены.
После второго запуска с SCD Type 1 целевая таблица отображает только текущее состояние:
| userId | city |
|---|---|
| 2 | Кармель |
| 3 | Лос-Анджелес |
| 4 | Долина смерти |
| 6 | Каньон Кингс |
Пример: Обработка моментальных снимков с помощью функций для управления версиями
Используйте этот подход, если вам нужен явный контроль над упорядочением моментальных снимков. Например, используйте этот подход при одновременном поступлении нескольких моментальных снимков или их прибытии в неправильном порядке. Напишите функцию, которая указывает, какой моментальный снимок обрабатывать следующим и его номер версии. API обрабатывает снэпшоты в порядке возрастания версии.
- Если несколько моментальных снимков находятся в хранилище, они обрабатываются по порядку.
- Если моментальный снимок поступает не по порядку (например,
snapshot_3поступает послеsnapshot_4), он пропускается. - Если новых моментальных снимков нет, функция возвращает
Noneи обработка не происходит.
Шаг 1. Подготовка файлов моментальных снимков
Создайте CSV-файлы, содержащие данные моментального снимка, и добавьте их в расположение тома или облачного хранилища. Назовите файлы в хронологическом порядке (например, snapshot_1.csv, snapshot_2.csv).
Каждый файл должен содержать столбцы для userId и city. Рассмотрим пример.
snapshot_1.csv:
| userId | city |
|---|---|
| 1 | Оахака |
| 2 | Монтеррей |
| 3 | Тихуана |
snapshot_2.csv:
| userId | city |
|---|---|
| 2 | Кармель |
| 3 | Лос-Анджелес |
| 4 | Долина смерти |
Шаг 2. Запуск AUTO CDC FROM SNAPSHOT с функцией версии
Создайте записную книжку и вставьте следующий код конвейера. Затем что такое конвейеры?.
from pyspark import pipelines as dp
from typing import Optional, Tuple
from pyspark.sql import DataFrame
def next_snapshot_and_version(latest_snapshot_version: Optional[int]) -> Optional[Tuple[DataFrame, int]]:
snapshot_dir = "/Volumes/main/cdc_tutorial/snapshots/" # or the location you created the sample data
files = dbutils.fs.ls(snapshot_dir)
snapshot_files = [f.name for f in files if f.name.startswith("snapshot_") and f.name.endswith(".csv")]
snapshot_versions = []
for filename in snapshot_files:
try:
version = int(filename.replace("snapshot_", "").replace(".csv", ""))
snapshot_versions.append(version)
except ValueError:
continue
snapshot_versions.sort()
if latest_snapshot_version is None:
if snapshot_versions:
next_version = snapshot_versions[0]
else:
return None
else:
next_versions = [v for v in snapshot_versions if v > latest_snapshot_version]
if next_versions:
next_version = next_versions[0]
else:
return None
snapshot_path = f"{snapshot_dir}snapshot_{next_version}.csv"
df = spark.read.format("csv").option("header", True).load(snapshot_path)
return (df, next_version)
dp.create_streaming_table("main.cdc_tutorial.target_versioned")
dp.create_auto_cdc_from_snapshot_flow(
target = "main.cdc_tutorial.target_versioned",
source = next_snapshot_and_version,
keys = ["userId"],
stored_as_scd_type = 2
)
Замечание
Чтобы вместо этого использовать SCD Type 1, задайте для него значение stored_as_scd_type=1.
После обработки snapshot_1.csvцелевая таблица содержит следующие записи:
| userId | city | __НАЧАТЬ_С | __END_AT |
|---|---|---|---|
| 1 | Оахака | 1 | ноль |
| 2 | Монтеррей | 1 | ноль |
| 3 | Тихуана | 1 | ноль |
После обработки snapshot_2.csvцелевая таблица содержит следующие записи:
| userId | city | __НАЧАТЬ_С | __END_AT |
|---|---|---|---|
| 1 | Оахака | 1 | 2 |
| 2 | Монтеррей | 1 | 2 |
| 2 | Кармель | 2 | ноль |
| 3 | Тихуана | 1 | 2 |
| 3 | Лос-Анджелес | 2 | ноль |
| 4 | Долина смерти | 2 | ноль |
Замечание
Помните, что для SCD Type 1 таблица выглядит точно так же, как последний моментальный снимок. Разница заключается в том, что последующие запросы могут использовать поток изменений для обработки только измененных записей.
Шаг 3. Добавление новых моментальных снимков
Добавьте новый CSV-файл в расположение хранилища с измененными данными (например, измененными значениями города, новыми строками или удаленными строками). Затем запустите конвейер снова, чтобы обработать новый снимок состояния.
Ограничения
- Столбец последовательности должен быть сортируемым типом данных.
NULLЗначения упорядочивания не поддерживаются. -
AUTO CDC FROM SNAPSHOTподдерживается только в интерфейсе конвейера Python; Интерфейс SQL не поддерживается. - Чтобы передавать данные в потоковом режиме из целевого объекта процесса AUTO CDC, читайте из его канала изменений. Дополнительные сведения см. в чтении потока изменяемых данных из целевой таблицы AUTO CDC.
Дополнительные ресурсы
- Отслеживание изменений данных и моментальные снимки: узнайте о понятиях CDC, моментальных снимках и типах SCD.
-
Репликация внешней таблицы RDBMS с помощью
AUTO CDC: узнайте, как выполнить начальную гидратацию с потокомonce, а затем продолжить обработку изменений. - Продвинутые темы AUTO CDC: Узнайте об операциях изменения в целевых объектах AUTO CDC, чтении потоков данных изменений и обработке метрик.
- Руководство. Создание конвейера ETL с помощью отслеживания измененных данных