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

Important

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

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

Колонка 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 столбец. Ключевые столбцы нельзя повторять, и их типы должны быть сортируемыми. Атомарные типы, такие как целые числа, строки и даты, могут быть ключами. MAP И VARIANT не могут быть ключами.

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

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

Note

Сведения о различиях в синтаксисе для автономных потоковых таблиц см. в разделе «Частичная замена моментального снимка с потоками REPLACE USING».

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 рассматривает строку, нарушающую условие, так, как будто источник никогда её не выдавал. Потерянная строка не заменяет, не удаляет и не изменяет соответствующие ключи в целевой таблице:

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

Поддерживаемые операции

В запросе потока поддерживаются следующие операции:

  • Проекция столбцов: выберите или перестройте подмножество столбцов.
  • Фильтры с WHERE.
  • Скалярные выражения, такие как CAST, арифметика и CASE.
  • Дедупликация с SELECT DISTINCT.
  • Пределы строк с LIMIT.
  • Функции-генераторы, такие как EXPLODE и POSEXPLODE.
  • Объединение двух потоковых источников с UNION ALL.
  • Соединения потока со статическими данными: внутреннее и левое внешнее, где поток находится с левой стороны.
  • Внутренние соединения между потоками.

Следующие операции требуют дополнительной конфигурации:

  • Окна, основанные на времени, такие как window(ts, '5 minutes'), требуют водяной метки.
  • Для внешних соединений поток-поток требуются водяная метка и условие по временному диапазону.

Ограничения

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

  • REPLACE USING поддерживает только один поток для каждой целевой таблицы. Комбинация REPLACE USING с другим типом потока для той же цели не поддерживается.
  • Целевая таблица должна быть создана в конвейере.

Следующие операции не поддерживаются в запросе потока:

  • Агрегации, такие как SUM, COUNT, и GROUP BY.
  • Оконные функции, кроме временных окон, например ROW_NUMBER() OVER (...), даже при наличии водяного знака.
  • Сортировка с ORDER BY.
  • Задавать операции, такие как INTERSECT и EXCEPT.
  • Профсоюзы, которые смешивают источник потока с не-стриминговым источником.
  • Чтение с нестримингового источника, например spark.range().
  • Полные и правые наружные соединения с статикой потока.

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")