API-интерфейсы AUTO CDC: упрощение отслеживания измененных данных с помощью конвейеров

Конвейеры 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.

Дополнительные ресурсы