Частичная замена снимка с использованием потоков REPLACE USING

Important

Эта функция доступна в бета-версии.

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

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

Как работает REPLACE USING

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

region_id тип_устройства event_type последовательный
1 iOS click 1
1 Android превращение 1
2 iOS click 1
2 настольный компьютер click 1

Поток REPLACE USING (region_id) SEQUENCE BY seq получает эти обновления для регионов 1 и 3. Во втором регионе обновлений нет:

region_id тип_устройства event_type последовательный
1 iOS click 2
1 Android превращение 2
1 настольный компьютер click 2
3 iOS click 1
3 настольный компьютер click 2

Цель становится:

region_id тип_устройства event_type последовательный Результат
1 iOS click 2 Заменён, потому что seq 2 больше seq 1
1 Android превращение 2 Заменён, потому что seq 2 больше seq 1
1 настольный компьютер click 2 Заменён, потому что seq 2 больше seq 1
2 iOS click 1 Не тронуты, потому что ключ отсутствует в этом обновлении
2 настольный компьютер click 1 Не тронуты, потому что ключ отсутствует в этом обновлении
3 настольный компьютер click 2 Добавлено. Строка seq 1 для региона 3 не добавляется, потому что применяется только самая высокая последовательность для ключа.

Требования

Потоки «Заменить с помощью» имеют следующие требования:

  • Потоки REPLACE USING выполняются в Databricks Runtime 18.2 и более поздних версий на классических или бессерверных вычислительных ресурсах. Databricks рекомендует Unity Catalog.
  • Источник должен быть источником потоковой передачи. REPLACE USING отклоняет источник без стриминга.
  • Вы должны указать как минимум один ключевой столбец и ровно один SEQUENCE BY столбец.

Когда использовать ЗАМЕНУ с использованием потоков

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

  • Используйте REPLACE USING, если ваш источник данных представляет собой серию частичных снимков, где ключом служит столбец. REPLACE USING перезаписывает только те данные, которые совпадают с входящими, оставляя все остальные данные нетронутыми. Для этого не требуется первичный ключ.
  • Используйте AUTO CDC, если ваш источник — это поток с записью данных изменений (CDC) с явными операциями вставки, обновления и удаления , или если вам нужна история медленно меняющегося размера (SCD) Type 2 . AUTO CDC также требует настоящего первичного ключа. См . API AUTO CDC: упрощение отслеживания изменений с помощью конвейеров.
  • Используйте REPLACE WHERE, если источник представляет собой снимок и вы хотите заново вычислить и перезаписать диапазон в целевой таблице, задаваемый предикатом, например данные за последние 7 дней, в пакетном режиме. Для этого не требуется первичный ключ. См. раздел "Пакетная обработка с помощью потоков REPLACEWHERE".

Создать ЗАМЕНУ С помощью потока

Определите REPLACE USING flows в SQL или Python.

SQL

Используйте FLOW REPLACE USINGусловие в одной строке с CREATE STREAMING TABLE:

CREATE STREAMING TABLE payments_current
FLOW REPLACE USING (payment_id) SEQUENCE BY payment_date BY NAME
SELECT payment_id, booking_id, status, payment_date
FROM STREAM(samples.wanderbricks.payments);

Кроме того, используйте синтаксис длинной формы CREATE FLOW :

CREATE STREAMING TABLE payments_current;

CREATE FLOW payments_flow AS
INSERT INTO payments_current BY NAME
REPLACE USING (payment_id) SEQUENCE BY payment_date
SELECT payment_id, booking_id, status, payment_date
FROM STREAM(samples.wanderbricks.payments);

Note

BY NAME требуется в SQL. Сопоставляет столбцы по именам, а не по их порядку.

Python

Объявим таблицу и поток вместе с @dp.table:

from pyspark import pipelines as dp

@dp.table(name="payments_current", replace_using=["payment_id"], sequence_by="payment_date")
def payments_current():
  return spark.readStream.table("samples.wanderbricks.payments")

Или укажите существующую потоковую таблицу с помощью @dp.replace_flow:

from pyspark import pipelines as dp

dp.create_streaming_table("payments_current")

@dp.replace_flow(target="payments_current", replace_using=["payment_id"], sequence_by="payment_date")
def payments_flow():
  return spark.readStream.table("samples.wanderbricks.payments")

replace_using — это список ключевых столбцов. sequence_by — это имя столбца или выражение Column и является обязательным, если задано replace_using.

Последовательность и неупорядоченные данные

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

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

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

Expectations

Потоки REPLACE USING поддерживают ожидания. warn И fail ведёт себя так же, как и в других потоках: warn продолжает нарушать строки, фиксирует нарушение и fail останавливает обновление. См. Управление качеством данных, используя ожидания конвейера.

Ожидание drop рассматривает строку, нарушающую условие, так, как будто источник никогда её не выдавал. Потерянная строка не заменяет, не удаляет и не изменяет соответствующие ключи в целевой таблице:

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

Ограничения

Потоки "REPLACE USING" имеют следующие ограничения:

  • REPLACE USING поддерживает только один поток для каждой целевой таблицы. Комбинация REPLACE USING с другим типом потока для той же цели не поддерживается.
  • Целевая таблица должна быть создана в конвейере.
  • Источник должен быть источником потоковой передачи.
  • Вы должны указать не менее одного ключевого столбца и столбец SEQUENCE BY. Ключевые столбцы нельзя повторять, и тип каждого ключевого столбца должен быть сортируемым. Атомарные типы, такие как целые числа, строки и даты, могут быть ключами, а MAPVARIANT не могут.
  • Сведения о различиях в синтаксисе для автономных потоковых таблиц см. в разделе «Частичная замена моментального снимка с потоками REPLACE USING».

Examples

Следующие примеры считывают данные из samples.wanderbricks.booking_updates, образца таблицы изменений статуса бронирования, доступного в каждом рабочем пространстве с поддержкой Unity Catalog. Каждое бронирование появляется один раз за смену, то есть booking_id повторяется с новым booking_update_id. См. набор данных Wanderbricks.

Пример 1: Сохраняйте последний запис для каждого ключа

Сохраняйте только текущее состояние каждого бронирования. Поток группируется по booking_id и упорядочивается по booking_update_id, так что самое последнее обновление для бронирования заменяет его более ранние обновления. Используйте AUTO CDC вместо этого, если ваш исходник — лента изменений с явными операциями вставки, обновления и удаления.

SQL

CREATE OR REFRESH STREAMING TABLE bookings_current
FLOW REPLACE USING (booking_id) SEQUENCE BY booking_update_id BY NAME
SELECT booking_id, status, total_amount, booking_update_id
FROM STREAM(samples.wanderbricks.booking_updates);

Python

from pyspark import pipelines as dp

@dp.table(
  name="bookings_current",
  replace_using=["booking_id"],
  sequence_by="booking_update_id"
)
def bookings_current():
  return spark.readStream.table("samples.wanderbricks.booking_updates")

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

Пример 2: Ключ по нескольким столбцам

Когда запись идентифицируется комбинацией столбцов, перечислите их все в REPLACE USING. Здесь каждое бронирование идентифицируется по (property_id, booking_id), поэтому поток сохраняет текущее состояние каждого бронирования на объект. Если ключевой столбец может содержать NULL, REPLACE USING сопоставляет NULL с NULL вместо того, чтобы пропускать строку.

SQL

CREATE OR REFRESH STREAMING TABLE bookings_by_property
FLOW REPLACE USING (property_id, booking_id) SEQUENCE BY booking_update_id BY NAME
SELECT property_id, booking_id, status, total_amount, booking_update_id
FROM STREAM(samples.wanderbricks.booking_updates);

Python

from pyspark import pipelines as dp

@dp.table(
  name="bookings_by_property",
  replace_using=["property_id", "booking_id"],
  sequence_by="booking_update_id"
)
def bookings_by_property():
  return spark.readStream.table("samples.wanderbricks.booking_updates")

Пример 3: Удаление недействительных записей с ожиданием

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

from pyspark import pipelines as dp

@dp.table(
  name="bookings_validated",
  replace_using=["booking_id"],
  sequence_by="booking_update_id"
)
@dp.expect_or_drop("positive_amount", "total_amount > 0")
def bookings_validated():
  return spark.readStream.table("samples.wanderbricks.booking_updates")