Примечание.
Для доступа к этой странице требуется авторизация. Вы можете попробовать войти или изменить каталоги.
Для доступа к этой странице требуется авторизация. Вы можете попробовать изменить каталоги.
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")