Репликация внешней таблицы RDBMS с помощью AUTO CDC

Вы можете реплицировать таблицу из внешней системы управления реляционными базами данных (RDBMS) в Azure Databricks с помощью AUTO CDC API в конвейерах. Вы узнаете:

  • Распространенные шаблоны настройки источников.
  • Как выполнить однократную полную копию существующих данных с помощью once потока.
  • Как непрерывно внедрять новые изменения посредством change потока.

Этот шаблон идеально подходит для создания медленно изменяющихся таблиц измерения (SCD) или синхронизации целевой таблицы с внешней системой записи.

Перед тем как начать

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

  • Полный снимок исходной таблицы в облачном хранилище. Этот набор данных используется для начальной загрузки.
  • Транзакционная лента изменений, размещенная в том же облачном хранилище (например, с помощью Debezium, Kafka или извлечения данных при изменении (CDC) на основе журналов). Этот поток данных является входными данными для текущего AUTO CDC процесса.

Настройка исходных представлений

Сначала определите два исходных представления для заполнения rdbms_orders целевой таблицы из пути orders_snapshot_pathк облачному хранилищу. Оба создаются в виде потоковых представлений по необработанным данным в облачном хранилище. Использование представлений обеспечивает более высокую эффективность, так как данные не должны быть записаны перед использованием в AUTO CDC процессе.

  • Первое исходное представление — это полный моментальный снимок (full_orders_snapshot)
  • Второй — это континуальный поток изменений (rdbms_orders_change_feed).

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

full_orders_snapshot()

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

Питон

Следующий пример Python:

  • Используется spark.readStream с автозагрузчиком (format("cloudFiles"))
  • Считывает JSON-файлы из каталога, определенного orders_snapshot_path
  • Устанавливает includeExistingFiles в true, чтобы обеспечить обработку исторических данных, уже находящихся на пути.
  • Задает inferColumnTypes для true автоматического определения схемы
  • Возвращает все столбцы с .select("\*")
@dp.view()
def full_orders_snapshot():
    return (
        spark.readStream
        .format("cloudFiles")
        .option("cloudFiles.format", "json")
        .option("cloudFiles.includeExistingFiles", "true")
        .option("cloudFiles.inferColumnTypes", "true")
        .load(orders_snapshot_path)
        .select("*")
    )

SQL

В следующем примере SQL параметры передаются в виде карты строковых пар "ключ-значение". orders_snapshot_path должен быть доступен как переменная SQL (например, определена с помощью параметров конвейера или интерполирована вручную).

CREATE OR REFRESH VIEW full_orders_snapshot
AS SELECT *
FROM STREAM read_files("${orders_snapshot_path}", "json", map(
  "cloudFiles.includeExistingFiles", "true",
  "cloudFiles.inferColumnTypes", "true"
));

rdbms_orders_change_feed()

На этом шаге создается второе представление, которое считывает добавочные данные об изменениях (например, из журналов CDC или таблиц изменений). Он считывает из orders_cdc_path и ожидает, что JSON-файлы в стиле CDC регулярно помещаются по этому пути.

Питон

@dp.view()
def rdbms_orders_change_feed():
return (
spark.readStream
.format("cloudFiles")
.option("cloudFiles.format", "json")
.option("cloudFiles.includeExistingFiles", "true")
.option("cloudFiles.inferColumnTypes", "true")
.load(orders_cdc_path)

SQL

В следующем примере SQL это переменная, которую можно интерполировать, ${orders_cdc_path} задав значение в параметрах конвейера или явно задав переменную в коде.

CREATE OR REFRESH VIEW rdbms_orders_change_feed
AS SELECT *
FROM STREAM read_files("${orders_cdc_path}", "json", map(
"cloudFiles.includeExistingFiles", "true",
"cloudFiles.inferColumnTypes", "true"
));

Начальная гидратация (единовременный поток)

Теперь, когда источники настроены, AUTO CDC логика объединяет оба источника в целевую потоковую таблицу. Во-первых, используйте однократный поток AUTO CDC с ONCE=TRUE для копирования полного содержимого таблицы RDBMS в потоковую таблицу. Это подготавливает целевую таблицу с историческими данными без ее повторного воспроизведения в будущих обновлениях.

Питон

from pyspark import pipelines as dp

# Step 1: Create the target streaming table

dp.create_streaming_table("rdbms_orders")

# Step 2: Once Flow — Load initial snapshot of full RDBMS table

dp.create_auto_cdc_flow(
  flow_name = "initial_load_orders",
  once = True,  # one-time load
  target = "rdbms_orders",
  source = "full_orders_snapshot",  # e.g., ingested from JDBC into bronze
  keys = ["order_id"],
  sequence_by = "timestamp",
  stored_as_scd_type = "1"
)

SQL


-- Step 1: Create the target streaming table
CREATE OR REFRESH STREAMING TABLE rdbms_orders;

-- Step 2: Once Flow for initial snapshot
CREATE FLOW rdbms_orders_hydrate
AS AUTO CDC ONCE INTO rdbms_orders
FROM stream(full_orders_snapshot)
KEYS (order_id)
SEQUENCE BY timestamp
STORED AS SCD TYPE 1;

Поток once выполняется только один раз. Новые файлы, которые добавлены после создания конвейера в full_orders_snapshot, игнорируются.

Это важно

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

Канал непрерывных изменений (поток изменений)

После начальной загрузки моментального снимка используйте другой AUTO CDC поток для непрерывной обработки изменений из потока данных CDC RDBMS. rdbms_orders Благодаря этому таблица обновляется с помощью вставок, обновлений и удалений.

Питон

from pyspark import pipelines as dp

# Step 3: Change Flow — Ingest ongoing CDC stream from source system

dp.create_auto_cdc_flow(
flow_name = "orders_incremental_cdc",
target = "rdbms_orders",
source = "rdbms_orders_change_feed", # e.g., ingested from Kafka or Debezium
keys = ["order_id"],
sequence_by = "timestamp",
stored_as_scd_type = "1"
)

SQL

-- Step 3: Continuous CDC ingestion
CREATE FLOW rdbms_orders_continuous
AS AUTO CDC INTO rdbms_orders
FROM stream(rdbms_orders_change_feed)
KEYS (order_id)
SEQUENCE BY timestamp
STORED AS SCD TYPE 1;

Соображения

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

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