Примечание.
Для доступа к этой странице требуется авторизация. Вы можете попробовать войти или изменить каталоги.
Для доступа к этой странице требуется авторизация. Вы можете попробовать изменить каталоги.
Потоки REPLACE WHERE в конвейерах Lakeflow пересчитывают и перезаписывают указанное подмножество таблицы без повторной обработки всей истории таблицы. Они обрабатывают данные, поступающие с задержкой, переобработку в вышестоящих системах, эволюцию схемы и дозагрузку исторических данных.
С помощью потока REPLACE WHERE вы определяете предикат в целевой таблице. Все строки, соответствующие предикату, удаляются и заменяются повторной оценкой исходного запроса для того же диапазона предиката. Строки, которые не соответствуют предикату, остаются без изменений.
Requirements
Потоки REPLACE WHERE имеют следующие требования:
- Databricks рекомендует каталог Unity и бессерверные вычисления. Добавочное обновление поддерживается только для бессерверных вычислений.
Когда следует использовать потоки REPLACE WHERE
Используйте потоки REPLACE WHERE для следующих сценариев:
- Добавочная пакетная обработка без семантики потоковой передачи. Обработка новых строк в пакетах без управления понятиями потоковой передачи, такими как подложки.
- Выборочная повторная обработка: пересчитывайте только те строки, которые соответствуют предикату, оставляя все остальные строки без изменений.
-
Сценарии, выходящие за рамки стандартных материализованных возможностей представления:
- Целевые таблицы с более длительным сроком хранения, чем у источника
- Предотвращение повторной компиляции при изменении таблицы измерений
- Эволюция схемы без повторной компиляции всей истории
Создание потока REPLACE WHERE
Определите потоки REPLACE WHERE в SQL или Python.
SQL
Используйте FLOW REPLACE WHEREусловие в одной строке с CREATE STREAMING TABLE:
CREATE STREAMING TABLE orders_enriched
FLOW REPLACE WHERE date >= date_add(current_date(), -7) BY NAME
SELECT
o.order_id,
o.date,
o.region,
p.product_name,
o.qty,
o.price
FROM orders_fct o
JOIN product_dim p
ON o.product_id = p.product_id;
Кроме того, используйте синтаксис длинной формы CREATE FLOW :
CREATE STREAMING TABLE orders_enriched;
CREATE FLOW orders_enriched AS
INSERT INTO orders_enriched BY NAME
REPLACE WHERE date >= date_add(current_date(), -7)
SELECT
o.order_id,
o.date,
o.region,
p.product_name,
o.qty,
o.price
FROM orders_fct o
JOIN product_dim p
ON o.product_id = p.product_id;
Python
В Python таблица и поток определяются в одной инструкции. Поток наследует то же имя, что и таблица:
from pyspark import pipelines as dp
from pyspark.sql import functions as F
from pyspark.sql.functions import col
@dp.table(
replace_where=col("date") >= F.date_sub(F.current_date(), 7)
)
def orders_enriched():
orders_fct = spark.read.table("orders_fct").select("date", "order_id", "region", "qty", "price")
product_dim = spark.read.table("product_dim")
return orders_fct.join(product_dim, "product_id")
Параметр replace_where принимает выражение столбца PySpark или строковый предикат.
В этих примерах все строки за последние 7 дней удаляются и orders_enriched повторно компилируются с помощью исходного запроса. Вам не нужно добавлять предикат в исходный запрос. Механизм конвейера автоматически применяет это при чтении данных из источника.
Note
BY NAME требуется в SQL. Сопоставляет столбцы по именам, а не по их порядку.
Заполнение исторических данных
Чтобы записать исторические или исправленные строки в целевую таблицу вне запланированных обновлений, выберите между двумя механизмами, в зависимости от того, где находятся исторические данные:
- Переопределения предиката: повторно выполните запрос-источник потока для разового диапазона предиката. Используется, когда исторические данные приходят из того же источника, что и добавочные данные.
- Инструкции DML: вставьте в целевую таблицу напрямую, обходя поток. Используется, если исторические данные находятся в другом источнике, чем добавочные данные.
Переопределения предикатов
Переопределите предикат REPLACE WHERE для одного обновления конвейера без изменения определения конвейера. Переопределения предиката являются однократными, применяются только к текущему обновлению и не влияют на будущие запуски.
Пример: начальная историческая загрузка
Чтобы выполнить однократное заполнение исторических данных при первой настройке конвейера:
pipeline_id = "<pipeline-id>"
overrides = [
{
"flow_name": "orders_enriched",
"predicate_override": "date BETWEEN '2020-01-01' AND '2024-12-31'",
}
]
resp = start_update_with_replace_where(
pipeline_id=pipeline_id,
replace_where_overrides=overrides,
)
print(resp)
Пример. Исправление столбца за определенный период
После обновления определения столбца примените это изменение к указанному историческому диапазону:
pipeline_id = "<pipeline-id>"
overrides = [
{
"flow_name": "orders_enriched",
"predicate_override": "date >= date_add(current_date(), -30)",
}
]
resp = start_update_with_replace_where(
pipeline_id=pipeline_id,
replace_where_overrides=overrides,
refresh_selection=["orders_enriched"],
)
print(resp)
Объедините несколько измерений в одном переопределяющем условии:
overrides = [
{
"flow_name": "orders_enriched",
"predicate_override": "date >= date_add(current_date(), -30) AND region = 'asia'",
}
]
Вспомогательные функции: start_update_with_replace_where
Используйте API обновления конвейера из записной книжки для отправки переопределений предиката:
from databricks.sdk import WorkspaceClient
from databricks.sdk.service.pipelines import StartUpdateResponse
def start_update_with_replace_where(
pipeline_id: str,
replace_where_overrides: list[dict],
refresh_selection: list[str] = None,
) -> StartUpdateResponse:
"""Start a pipeline update with REPLACE WHERE predicate overrides."""
client = WorkspaceClient()
body = {
"pipeline_id": pipeline_id,
"cause": "JOB_TASK",
"update_cause_details": {
"job_details": {"performance_target": "PERFORMANCE"}
},
"replace_where_overrides": replace_where_overrides,
}
if refresh_selection:
body["refresh_selection"] = refresh_selection
res = client.api_client.do(
"POST",
f"/api/2.0/pipelines/{pipeline_id}/updates",
body=body,
headers={"Accept": "application/json", "Content-Type": "application/json"},
)
return StartUpdateResponse.from_dict(res)
Инструкции DML
Выполняйте инструкции DML непосредственно в целевой таблице вне конвейера, чтобы выполнить первоначальную загрузку или корректировки, например загрузку данных из устаревшей таблицы:
INSERT INTO orders_enriched
SELECT *
FROM orders_enriched_legacy
WHERE date < '2025-01-01';
Строки, вставленные через DML, не подпадают под предикат REPLACE WHERE и сохраняются при плановых обновлениях, если только они не попадают в диапазон предиката при будущем запуске.
Поведение полного обновления
Полное обновление потока REPLACE WHERE повторно выполняет исходный запрос с использованием только текущего предиката. Строки, вставленные в результате переопределений предикатов или выполнения операторов DML вне текущего диапазона предиката, удаляются безвозвратно.
Предупреждение
Полное обновление очищает все существующие данные и повторно выполняет поток с помощью только определенного предиката. Если конвейер работает в течение года с предикатом 7 дней, полное обновление приводит к тому, что таблица содержит только последние 7 дней данных. Все старые строки окончательно удаляются.
Чтобы предотвратить полное обновление таблицы, задайте для свойства таблицы значение pipelines.reset.allowedfalse. См. справочник по свойствам конвейера.
Добавочное обновление
Потоки REPLACE WHERE используют инкрементное обновление, когда это возможно, обрабатывая только те исходные данные, которые изменились с момента последнего обновления, вместо пересчёта всего окна замены целиком. Добавочное обновление требует бессерверных вычислений.
Когда применяется инкрементальное обновление
Все из перечисленного ниже должно быть истинным:
- Конвейер выполняется на бессерверных вычислениях.
- Поддерживается структура запроса. Сведения о поддерживаемом наборе операторов см. в инкрементальном обновлении.
- Предикат ссылается на базовые столбцы из исходной таблицы. Предикаты по производным значениям, таким как результаты агрегатных или оконных функций, нельзя перенести на источник, из-за чего становится недоступным инкрементное обновление.
- Внешний DML не изменил строки в текущем окне замены. DML, изменяющий строки за пределами текущего окна, не затрагивается.
- Текущее окно замены не включает строки, исключенные из предыдущего предиката. Если вы расширяете предикат так, чтобы он охватывал диапазон, который ранее не обрабатывался, это обновление выполняется с полным пересчётом. Последующие обновления имеют право на добавочное обновление снова.
- Предикат детерминирован. Предикаты, использующие недетерминированные функции, такие как
rand(), отключают инкрементное обновление. Темпоральные функции, такие какcurrent_date()разрешены.
Первое обновление любого потока всегда является полным вычислением. Если какое-либо условие не выполнено, это обновление возвращается к полной повторной компиляции текущего окна замены.
Рекомендации по инкрементальному обновлению
Следуйте этим рекомендациям, чтобы потоки REPLACE WHERE оставались допустимыми для добавочного обновления.
Используйте подвижную нижнюю границу
Предикаты с перемещением нижней границы остаются допустимыми для добавочного обновления на неопределенный срок.
FLOW REPLACE WHERE date >= date_add(current_date(), -7)
Перемещаемая верхняя граница, например date BETWEEN date_add(current_date(), -7) AND current_date(), может сдвинуть окно, чтобы включить ранее исключенные строки, активировав обратный возврат к полной повторной компиляции.
Включить столбец предикатов в GROUP BY
При агрегации включите столбец условия в GROUP BY, чтобы механизм мог протолкнуть условие ниже операции агрегации.
FLOW REPLACE WHERE date >= date_add(current_date(), -7) BY NAME
SELECT date, region, SUM(amount) AS total
FROM sales
GROUP BY date, region;
Если столбец предиката отсутствует в GROUP BY, предикат нельзя протолкнуть ниже агрегации, и источник сканируется полностью.
Включить столбец предиката в ключи соединения
Включите столбец предиката в условие соединения, чтобы механизм мог отсечь все соединяемые источники.
FLOW REPLACE WHERE date >= date_add(current_date(), -7) BY NAME
SELECT f.date, f.user_id, d.region, f.revenue
FROM fact f
JOIN dim d ON f.date = d.date AND f.user_id = d.user_id;
Если присоединённая таблица не содержит столбец предиката, при каждом обновлении она полностью сканируется.
Диагностировать переход к полному пересчёту
Когда обновление переходит к полному пересчёту, причина указывается в событии planning_information для потока. См. статью "Мониторинг журналов событий конвейера". В следующей таблице перечислены причины, сообщаемые в событии:
| Причина | Meaning |
|---|---|
EXTERNAL_CHANGE_IN_REPLACE_WINDOW |
Внешняя операция DML изменила строки в текущем окне замены. |
REPLACE_WHERE_NOT_DETERMINISTIC |
Предикат использует недетерминированные выражения. |
PRIOR_REPLACE_WHERE_NOT_DETERMINISTIC |
Предыдущее обновление использовало недетерминированный предикат. |
UNSUPPORTED_REPLACE_WHERE_PREDICATE |
Предикат нельзя отправить в любой источник, текущее окно содержит строки, не обработанные предыдущим предикатом, или выполнение использует переопределение предиката. |
Ограничения
Потоки REPLACE WHERE имеют следующие ограничения:
- Целевая таблица должна быть создана в конвейере.
- Для целевой таблицы допускается только один поток REPLACE WHERE .
- Таблица, используемая потоком REPLACE WHERE в качестве целевой, не может одновременно использоваться и другим типом потока, например потоком AUTO CDC или потоком дозаписи.
- Ожидаемые значения не поддерживаются в таблицах, в которые записывают потоки REPLACE WHERE.
- Описание синтаксиса и различий при обратном заполнении для автономных потоковых таблиц см. в разделе REPLACE WHERE потоки для автономных потоковых таблиц.
Примеры
В следующих примерах показаны распространенные шаблоны потока REPLACE WHERE .
Пример 1. Сохранение статистических статистических данных из источника ограниченного хранения
Этот пример сохраняет ежедневные статистические данные неограниченное время даже после того, как необработанные данные выходят из исходной таблицы (3-дневное хранение):
SQL
CREATE STREAMING TABLE events_agg
FLOW REPLACE WHERE date >= date_add(current_date(), -3) BY NAME
SELECT
date,
key,
SUM(val) AS agg
FROM events_raw
GROUP BY ALL;
Python
from pyspark import pipelines as dp
from pyspark.sql import functions as F
from pyspark.sql.functions import col
@dp.table(
replace_where=col("date") >= F.date_sub(F.current_date(), 3)
)
def events_agg():
return (
spark.read.table("events_raw")
.groupBy("date", "key")
.agg(F.sum("val").alias("agg"))
)
Пример 2. Предотвращение повторной компиляции при изменении таблицы измерений
В этом примере строки фактов сохраняются без изменений при изменении атрибутов измерения:
SQL
CREATE STREAMING TABLE fact_dim_join
FLOW REPLACE WHERE f.date >= date_add(current_date(), -1) BY NAME
SELECT
f.date,
f.user_id,
d.region,
f.revenue
FROM fact_table f
JOIN dim_users d
ON f.user_id = d.user_id;
Python
from pyspark import pipelines as dp
from pyspark.sql import functions as F
from pyspark.sql.functions import col
@dp.table(
replace_where=col("date") >= F.date_sub(F.current_date(), 1)
)
def fact_dim_join():
fact_table = spark.read.table("fact_table").alias("f")
dim_users = spark.read.table("dim_users").alias("d")
return (
fact_table.join(dim_users, col("f.user_id") == col("d.user_id"))
.select(
col("f.date"),
col("f.user_id"),
col("d.region"),
col("f.revenue"),
)
)
Если регион пользователя изменяется, перекомпьютируются только последние строки. Исторические строки сохраняют значение региона во время их записи. Чтобы исправить исторические записи, выполните целевой бэкфилл с помощью переопределений предикатов.
Пример 3. Добавление новой метрики без повторной компиляции полной истории
В этом примере показано, как изменить определение таблицы и заполучить только целевой диапазон:
Определите начальную таблицу:
SQL
CREATE STREAMING TABLE clickstream_daily FLOW REPLACE WHERE event_date >= date_add(current_date(), -7) BY NAME SELECT event_date, page_id, COUNT(*) AS clicks FROM clickstream_raw GROUP BY ALL;Python
from pyspark import pipelines as dp from pyspark.sql import functions as F from pyspark.sql.functions import col @dp.table( replace_where=col("event_date") >= F.date_sub(F.current_date(), 7) ) def clickstream_daily(): return ( spark.read.table("clickstream_raw") .groupBy("event_date", "page_id") .agg(F.count("*").alias("clicks")) )Обновите запрос, чтобы добавить
uniq_users:SQL
CREATE STREAMING TABLE clickstream_daily FLOW REPLACE WHERE event_date >= date_add(current_date(), -7) BY NAME SELECT event_date, page_id, COUNT(*) AS clicks, COUNT(DISTINCT user_id) AS uniq_users FROM clickstream_raw GROUP BY ALL;Python
@dp.table( replace_where=col("event_date") >= F.date_sub(F.current_date(), 7) ) def clickstream_daily(): return ( spark.read.table("clickstream_raw") .groupBy("event_date", "page_id") .agg( F.count("*").alias("clicks"), F.countDistinct("user_id").alias("uniq_users"), ) )Дозаполнить новую метрику за последние 30 дней:
overrides = [ { "flow_name": "clickstream_daily", "predicate_override": "event_date BETWEEN '2026-01-01' AND '2026-01-30'", } ] resp = start_update_with_replace_where( pipeline_id="<pipeline-id>", replace_where_overrides=overrides, refresh_selection=["clickstream_daily"], )Строки, которые старше диапазона с дозаполнением, содержат
NULLдляuniq_users.
Пример 4: Выполните итерации на небольшом временном диапазоне перед дозаполнением всей истории
В этом примере показано, как проверить логику запроса в небольшом окне данных перед обработкой полного исторического диапазона.
Начните с небольшого окна, чтобы при каждом обновлении пересчитывались только последние 7 дней, пока вы пересматриваете запрос:
SQL
CREATE STREAMING TABLE revenue_attribution
FLOW REPLACE WHERE event_date >= date_add(current_date(), -7) BY NAME
SELECT
event_date,
campaign_id,
SUM(revenue) AS total_revenue
FROM marketing_events
GROUP BY ALL;
Python
from pyspark import pipelines as dp
from pyspark.sql import functions as F
from pyspark.sql.functions import col
@dp.table(
replace_where=col("event_date") >= F.date_sub(F.current_date(), 7)
)
def revenue_attribution():
return (
spark.read.table("marketing_events")
.groupBy("event_date", "campaign_id")
.agg(F.sum("revenue").alias("total_revenue"))
)
После окончательной настройки запроса используйте переопределение предиката, чтобы выполнить однократную историческую догрузку данных:
overrides = [
{
"flow_name": "revenue_attribution",
"predicate_override": "event_date >= date_add(current_date(), -365)",
}
]
resp = start_update_with_replace_where(
pipeline_id="<pipeline-id>",
replace_where_overrides=overrides,
refresh_selection=["revenue_attribution"],
)