Поток изменений данных Lakebase

Note

Функция Lakebase Change Data Feed находится в общедоступной предварительной версии.

Что такое поток данных об изменениях Lakebase?

Lakebase представляет встроенный поток изменений данных (CDF), открывая доступ к вашим операционным данным для последующих конвейеров обработки, моделей и приложений. Каждая операция вставки, обновления и удаления в таблице Lakebase Postgres считывается из журнала упреждающей записи (WAL) и сохраняется как новая строка в управляемой Delta-таблице Unity Catalog; записи группируются в пакеты и записываются примерно каждые 15 секунд. Журнал изменений хранится в открытом формате, считываемом любым вычислительным ядром.

Целевые таблицы соответствуют той же форме, что и канал данных Delta Change: каждая строка содержит _pg_change_typeномер LSN, идентификатор транзакции и метку времени. Операционные изменения становятся источником первого класса для ETL, аудита и подчиненных потребителей — без создания внешнего стека CDC.

Поток данных Lakebase CDF из Postgres через wal2delta в таблицы Delta в каталоге Unity.

Случаи использования

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

Сценарий использования Description
Конвейеры ETL Используйте Lakebase в качестве бронзового источника для конвейеров медальона. Создайте инкрементные конвейеры Lakeflow или задания Spark Structured Streaming на основе потока изменений и обновите нижестоящие таблицы уровня Silver и Gold.
Журналы аудита Сохраняйте полную, запрашиваемую историю каждой вставки, обновления и удаления в таблице Lakebase для соответствия требованиям и судебной экспертизы. История неизменяема.
Внешние системы Сохраните данные об изменении Lakebase в открытом формате, который может использовать любой обработчик. Поскольку целевым объектом является таблица Delta в Unity Catalog, внешние системы и читатели, не относящиеся к Databricks, могут получать доступ к этому потоку данных напрямую.

Включение этой предварительной версии

Администратор рабочей области должен включить предварительную версию Lakebase Change Data Feed на странице Предварительные версии рабочей области.

Requirements

  • Лейкбейс:Проект Lakebase, работающий под управлением Postgres 16, 17 или 18.
  • Исходная база данных: Таблицы исходников могут находиться в любой отдельной базе данных вашего проекта Lakebase. CDF фиксирует изменения из одной базы данных на каждый поток; Это не ограничивается базой databricks_postgres данных, с которой создаётся каждый проект.
  • Каталог Unity: Для учетной записи, настраивающей USE CATALOGCDF, требуются USE SCHEMA и CREATE TABLE для целевых каталога и схемы. См . раздел "Предоставление разрешений" для объекта.
  • Хранилище по умолчанию: Конечные каталоги, настроенные с хранилищем по умолчанию, не поддерживаются.
  • Проект Lakebase: Для роли Postgres требуются разрешения CAN MANAGE в проекте Lakebase. Владельцы проекта по умолчанию МОГУТ УПРАВЛЯТЬ. См. раздел "Управление разрешениями проекта".
  • Типы данных: См. сопоставление типов данных. Типы без прямого эквивалента Delta хранятся в виде STRING.

Note

Free Edition:Рабочие пространства Databricks Free Edition используют стандартное хранилище для каталога, созданного с этим рабочим пространством. Чтобы использовать Lakebase CDF, создайте каталог , управляемое место хранения которого является внешним местоположением.

Настройте Lakebase CDF

Чтобы приступить к работе, установите для таблиц, которые должны входить в поток, параметр REPLICA IDENTITY FULL (Шаг 1), а затем запустите CDF в приложении Lakebase (Шаг 2). Ваши данные отображаются как таблицы Delta lb_<table_name>_history в каталоге и схеме Unity Catalog, которые вы выберете.

Note

Вы можете запустить CDF из интерфейса Lakebase или с помощью API. Для программного управления лентой используйте операции CDF в REST API Postgres и SDK Databricks для создания ленты, проверки её состояния, отключения или удаления конфигурации. См. раздел «Change Data Feed » в руководстве Lakebase API.

Шаг 1. Установите REPLICA IDENTITY FULL

Чтобы таблица Lakebase участвовала в CDF, для неё должен быть установлен параметр REPLICA IDENTITY FULL. По умолчанию Postgres регистрирует только первичный ключ при обновлении или удалении строки. Установка полной идентичности предписывает Postgres записывать в журнал предзаписи (WAL) состояние строки до и после изменения, что необходимо CDF для построения полной истории изменений.

Эти команды можно выполнить в редакторе SQL Lakebase или любом клиенте Postgres.

Одна таблица

ALTER TABLE <table_name> REPLICA IDENTITY FULL;

Все существующие таблицы в схеме

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

DO $$
DECLARE r record;
BEGIN
  FOR r IN
    SELECT table_schema, table_name
    FROM information_schema.tables
    WHERE table_schema = 'public'
      AND table_type = 'BASE TABLE'
  LOOP
    EXECUTE format(
      'ALTER TABLE %I.%I REPLICA IDENTITY FULL;',
      r.table_schema, r.table_name
    );
  END LOOP;
END $$;

Автоматическое применение к будущим таблицам

Чтобы каждая вновь создаваемая таблица автоматически получала REPLICA IDENTITY FULL, установите триггер события PostgreSQL. Он выполняется после каждого CREATE TABLE и задает свойство IDENTITY для новой таблицы:

CREATE OR REPLACE FUNCTION public.set_full_replica_identity()
RETURNS event_trigger
LANGUAGE plpgsql
AS $$
DECLARE
  obj record;
BEGIN
  FOR obj IN
    SELECT * FROM pg_event_trigger_ddl_commands()
    WHERE command_tag = 'CREATE TABLE'
  LOOP
    EXECUTE format(
      'ALTER TABLE %s REPLICA IDENTITY FULL;',
      obj.object_identity
    );
  END LOOP;
END $$;

CREATE EVENT TRIGGER set_full_replica_identity_on_create
ON ddl_command_end
WHEN TAG IN ('CREATE TABLE')
EXECUTE FUNCTION public.set_full_replica_identity();

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

Проверьте, для каких таблиц задан идентификатор реплики

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

SELECT n.nspname AS table_schema,
       c.relname AS table_name,
       CASE c.relreplident
         WHEN 'd' THEN 'default'
         WHEN 'n' THEN 'nothing'
         WHEN 'f' THEN 'full'
         WHEN 'i' THEN 'index'
       END AS replica_identity
FROM pg_class c
JOIN pg_namespace n ON n.oid = c.relnamespace
WHERE c.relkind = 'r'
  AND n.nspname = 'public'
ORDER BY n.nspname, c.relname;

Только строки с replica_identity = 'full' готовы для CDF.

API:REPLICA IDENTITY FULL — это стандартный DDL Postgres. См. ссылку на PostgreSQLALTER TABLE.

Шаг 2: Запустите канал передачи данных об изменениях

Lakebase CDF конфигурируется на уровне схемы. После запуска все текущие и будущие таблицы в исходной схеме включаются в поток.

  1. В рабочей области Azure Databricks откройте Lakebase Postgres из переключателя приложений (в правом верхнем углу).
  2. Выберите проект Lakebase и ветвь, которую вы хотите использовать (например, рабочую или основную).
  3. Откройте Обзор ветви, щелкнув по имени ветви в верхней цепочке навигации, затем нажмите вкладку Lakebase CDF.
  4. Нажмите Пуск.
  5. В диалоговом окне конфигурации:
    • База данных: Выберите исходную базу данных Postgres. Вы можете выбрать любую базу данных в проекте, и обязательно должны выбрать одну, даже если в проекте только одна база данных.
    • Схемы: Выберите исходную схему Postgres.
    • В каталог: Выберите целевой каталог в Unity Catalog.
    • Схема: Выберите целевую схему каталога Unity.
  6. Нажмите Start, чтобы запустить поток.

Обзор ветви с вкладкой Lakebase CDF, где показаны Start и настройки схемы.

Таблицы отображаются в назначении как lb_<table_name>_history. Чтобы найти их, откройте каталог на боковой панели, перейдите к целевому каталогу и схеме и откройте вкладку "Таблицы ".

Вкладка Lakebase CDF содержит две подвкладки:

Подвкладки показывают сопоставление и прогресс по каждой таблице.

  • Схемы: Перечисляет каждую исходную схему, ее целевой каталог и схему в каталоге Unity и состояние.
  • Таблицы: Перечислены все исходные таблицы, соответствующая целевая таблица lb_<table_name>_history, статус (Streaming или Snapshotting), Зафиксированный LSN (насколько далеко поток записал данные в Delta; пока идет создание исходного снимка, отображается как -), а также Последнее обновление (время последнего получения таблицей изменений).

Вы также можете проверить состояние канала из Postgres, выполнив это в редакторе SQL Lakebase:

SELECT * FROM wal2delta.tables;

Результат включает table_oid, status (STREAMING или SNAPSHOTTING), committed_lsn и last_write_time согласно таблице.

Important

Что такое wal2delta? Lakebase CDF работает с расширением wal2delta Postgres, которое выполняется внутри вычислительных ресурсов Lakebase. Он использует логическое декодирование для извлечения изменений из журнала предзаписи (WAL) и записывает их в таблицы Delta в Unity Catalog.

API: Чтобы программно получить конфигурацию потока и статус по таблицам, см. раздел «Изменить поток данных » в руководстве Lakebase API.

Схема целевой таблицы

CDF записывает одну таблицу Delta для каждой исходной таблицы с именем lb_<table_name>_history в целевых каталоге и схеме. Помимо исходных столбцов каждая строка содержит следующие системные столбцы:

Column Тип Description
_pg_change_type ТЕКСТ Тип операции: insert, delete, update_preimageили update_postimage.
_pg_lsn BIGINT Номер последовательности журнала Postgres.
_pg_xid INTEGER Идентификатор транзакции Postgres.
_timestamp TIMESTAMP Метка времени обработки изменения (без часового пояса).
_sort_by BIGINT Монотонный ключ сортировки, используемый для упорядочивания всех изменений.

Распространенные шаблоны изменений

  • Начальный снимок: При первом запуске CDF для существующей таблицы Lakebase каждая существующая строка записывается с помощью _pg_change_type = 'insert'.
  • Обновления: Обновление создает две строки: один с _pg_change_type = 'update_preimage' (старой строкой) и один с _pg_change_type = 'update_postimage' (новой строкой).
  • Удаляет: Удаление создает одну строку с _pg_change_type = 'delete'.

Это те же события изменений, что и в Delta Change Data Feed, поэтому применяются те же шаблоны последующей обработки.

Оперативное поведение

  • Конфликты имён: Если две исходные таблицы сопоставляются с одним и тем же именем назначения (например, sales.users и marketing.users обе сопоставляются с lb_users_history), CDF записывает первую в lb_users_history и автоматически добавляет суффикс ко второй, получая lb_users_history_1. Вы можете переименовать любую из двух целевых таблиц в Unity Catalog, и поток продолжит работать.
  • Область уровня схемы: При запуске CDF на схеме Lakebase будет включена каждая текущая и будущая таблица в этой схеме. Пустые таблицы пропускаются — таблица должна содержать хотя бы одну строку, чтобы появиться в целевом месте.
  • Удаленные исходные таблицы: При удалении таблицы в Lakebase сохраняется целевая таблица Delta в каталоге Unity.

Создание подчиненных конвейеров

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

Пример сценария. Приложение для электронной коммерции записывает заказы в таблицу Postgres orders, при этом каждая строка содержит item_id и quantity. Логистической команде нужны актуальные данные об уровне складских запасов. С CDF каждое изменение в orders сохраняется в таблице Delta lb_orders_history в Unity Catalog. Последующие конвейеры обработки считывают этот поток изменений и обновляют таблицу inventory_levels всякий раз, когда заказ оформляется, изменяется или отменяется.

Вычисление текущей инвентаризации с материализованным представлением

Самый простой подход — SQL-материализованное представление на основе таблицы истории. MV инкрементально обновляется по мере поступления новых событий об изменениях, а последующие потребители обращаются к нему с запросами как к любой другой таблице.

CREATE MATERIALIZED VIEW inventory_levels AS
SELECT
  item_id,
  SUM(
    CASE
      -- New orders (and the "new half" of updates) decrement inventory
      WHEN _pg_change_type IN ('insert', 'update_postimage') THEN -quantity
      -- Cancellations (and the "old half" of updates) restore inventory
      WHEN _pg_change_type IN ('delete', 'update_preimage') THEN quantity
      ELSE 0
    END
  ) AS current_inventory,
  MAX(_timestamp) AS last_transaction_ts,
  MAX(_pg_lsn) AS last_lsn
FROM lb_orders_history
GROUP BY item_id;

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

Потоковая обработка изменений с помощью декларативных конвейеров Spark

Для структурированной архитектуры медальона используйте конвейеры Lakeflow для объявления бронзовых, серебряных и золотых таблиц. Конвейеры Lakeflow запускают их как единый связанный конвейер, автоматически выполняя контрольные точки и управление зависимостями.

import dlt
from pyspark.sql import functions as F

@dlt.table
def inventory_adjustments():
    return (
        spark.readStream.table("<catalog>.<schema>.lb_orders_history")
        .withColumn(
            "delta",
            F.when(F.col("_pg_change_type").isin("insert", "update_postimage"), -F.col("quantity"))
             .when(F.col("_pg_change_type").isin("delete", "update_preimage"), F.col("quantity"))
             .otherwise(0),
        )
        .select("item_id", "delta", "_timestamp")
    )

@dlt.expect_or_drop("non_negative_stock", "on_hand >= 0")
@dlt.table
def inventory_levels():
    return (
        spark.read.table("LIVE.inventory_adjustments")
        .groupBy("item_id")
        .agg(F.sum("delta").alias("on_hand"))
    )

inventory_adjustments инкрементально считывает lb_orders_history с помощью readStream и формирует дельту для каждого события. inventory_levels выполняет агрегацию по item_id для вычисления текущего запаса. Проверка отбрасывает строки, которые привели бы к отрицательному остатку, что указывает на ошибку на более раннем этапе обработки.

Полное пошаговое руководство см. в руководстве по созданию конвейера ETL с помощью отслеживания измененных данных.

Настраиваемая обработка с помощью структурированной потоковой передачи Spark

Если вам нужен полный контроль — например, для пользовательских слияний, побочных эффектов или нескольких приемников, — читайте таблицу истории напрямую с помощью Spark Structured Streaming и используйте foreachBatch для записи в целевое расположение.

from pyspark.sql import functions as F
from delta.tables import DeltaTable

def update_inventory(batch_df, batch_id):
    deltas = (
        batch_df
        .withColumn(
            "delta",
            F.when(F.col("_pg_change_type").isin("insert", "update_postimage"), -F.col("quantity"))
             .when(F.col("_pg_change_type").isin("delete", "update_preimage"), F.col("quantity"))
             .otherwise(0),
        )
        .groupBy("item_id")
        .agg(F.sum("delta").alias("delta"))
    )

    target = DeltaTable.forName(spark, "<catalog>.<schema>.inventory_levels")
    (target.alias("t")
        .merge(deltas.alias("s"), "t.item_id = s.item_id")
        .whenMatchedUpdate(set={"on_hand": F.expr("t.on_hand + s.delta")})
        .whenNotMatchedInsert(values={"item_id": "s.item_id", "on_hand": "s.delta"})
        .execute())

(spark.readStream.table("<catalog>.<schema>.lb_orders_history")
    .writeStream
    .foreachBatch(update_inventory)
    .option("checkpointLocation", "/Volumes/<catalog>/<schema>/checkpoints/inventory_levels")
    .start())

Каждый микропакет агрегирует события изменений по item_id и объединяет итоговые дельты в inventory_levels.

Изначально рассчитан на поэтапное развитие. Каждая таблица lb_<table_name>_history — это таблица Delta только с добавлением данных. Каждое изменение источника записывается как новая строка с _pg_change_type маркировкой операции. Databricks SQL материализованные представления, потоки в конвейерах Lakeflow и задания Spark Structured Streaming все обрабатывают новые строки инкрементально из журнала транзакций Delta, поэтому последующие конвейеры выполняют только объём работы, пропорциональный внесённым изменениям. Вам не нужно включать Delta Change Data Feed для таблицы журнала, поскольку семантика изменений уже закодирована в данных строк.

Сопоставление типов данных

CDF поддерживает большинство стандартных типов примитивов PostgreSQL. Типы без прямого эквивалента Delta хранятся в виде STRING.

Тип PostgreSQL Тип Delta в Azure Databricks Примечания.
BOOLEAN BOOLEAN
INT, SMALLINT, BIGINT INT, SMALLINT, BIGINT
ТЕКСТ, VARCHAR, CHAR СТРУНА
JSONB СТРУНА Хранится как строка JSON.
ENUM СТРУНА Хранится как метка перечисления.
ЧИСЛОВОЙ / ДЕСЯТИЧНЫЙ ДЕСЯТИЧНОЕ ЧИСЛО или СТРОКА По возможности использует точность и разрядность исходного значения. Выполняет изменение масштаба без потери данных для несовместимых значений точности и масштаба. Возвращается к STRING, когда точность превышает 38 или когда точность или масштабирование не определены (необязанные ЧИСЛОВЫЕ значения). Все столбцы NUMERIC/DECIMAL могут иметь значение NULL, так как значения NaN сопоставляются со значением NULL. См. раздел "Числовые типы PostgreSQL".
DATE DATE
TIMESTAMP TIMESTAMP_NTZ
TIMESTAMPTZ TIMESTAMP
ПЛАВАЙТЕ, УДВАЯ ПЛАВАЙТЕ, УДВАЯ

Типы, хранящиеся в виде STRING:

  • Geography/Geometry (PostGIS): Типы из расширения PostGIS (например, geometry, geography).
  • Vector (pgvector): Тип vector из расширения pgvector.
  • Составные/структурные типы: Пользовательские типы, определяемые с помощью CREATE TYPE ... AS (field_name type, ...). Это типы строк, похожие на строки с именованными полями.
  • Карта: Типы ключей, такие как hstore (из hstore расширения). Postgres не имеет встроенного типа карты. hstore — это распространенный способ хранения пар "ключ-значение" в столбце.

Управление изменениями схемы

  • Переименование таблицы в Postgres (например, ALTER TABLE users RENAME TO customers) позволяет веб-каналу продолжать работу. Имя целевой таблицы Delta не изменится — остается lb_users_history.
  • Изменения схемы (добавление столбца, удаление столбца или изменение типа данных столбца) запускают повторное создание снимка затронутой таблицы. CDF повторно считывает всю таблицу из Postgres и перезаписывает ее в целевую таблицу Delta.

Отключить Lakebase CDF

Отключение CDF прекращает поток данных для всех схем Lakebase в проекте.

  1. В рабочей области Azure Databricks откройте Lakebase Postgres из переключателя приложений (в правом верхнем углу).
  2. Выберите проект Lakebase и ветвь, в которой вы настроили CDF.
  3. Откройте Обзор ветви, щелкнув по имени ветви в верхней цепочке навигации, затем нажмите вкладку Lakebase CDF.
  4. Нажмите кнопку "Отключить". В диалоговом окне подтверждения просмотрите предупреждение о том, что изменения перестают передаваться в таблицы Delta, а затем нажмите кнопку "Отключить ", чтобы подтвердить.

Отключение CDF не перезапускает вычислительные ресурсы.

API: Чтобы программно отключить или удалить конфигурацию ленты, см. раздел «Изменить поток данных » в руководстве Lakebase API.

Ограничения и устранение неполадок

Вы можете увидеть статус каждой таблицы (создание снимка, пропущена или потоковая обработка) на вкладке Lakebase CDF или выполнив это в Lakebase:

SELECT * FROM wal2delta.tables;

Распространённые причины, по которым таблица не отображается в ленте:

  • REPLICA IDENTITY FULL не задано: выполните команду ALTER TABLE <table_name> REPLICA IDENTITY FULL; для таблицы. См. Шаг 1: Установите полную идентификацию реплики.
  • Секционированные таблицы: Секционированные таблицы Lakebase не поддерживаются. Схема, содержащая секционированные таблицы, приводит к сбою этих таблиц.
  • Пустые таблицы: Таблица с нулевыми строками пропускается до тех пор, пока не существует хотя бы одна строка.

Предупреждение

Не изменяйте целевые lb_<table_name>_history таблицы Delta следующими способами:

Note

Приватная конечная точка на хранилище назначения: Lakebase CDF не поддерживается, когда управляемое хранилище для каталога Unity Catalog доступно только через приватную конечную точку. Примеры включают конечную точку интерфейса AWS PrivateLink или приватную конечную точку Azure с отключённым доступом к публичной сети к учетной записи хранилища. В качестве обходного пути настройте каталог, чьё управляемое хранилище доступно публично, и используйте этот каталог в качестве назначения для CDF.

Дальнейшие действия