REPLACE WHERE 流程會重新計算並覆寫獨立串流資料表的指定子集,而無需重新處理整個資料表的歷史資料。 它們處理遲到的資料、上游重新處理、結構演化及回填。
在 REPLACE WHERE 流程中,您會在目標資料表上定義一個條件述詞。 所有與謂詞相符的列都會被刪除,並以重新評估該謂詞範圍的來源查詢來取代。 不符合述詞條件的列會維持不變。
Requirements
REPLACE WHERE 流程具備以下需求:
- Databricks 推薦使用 Unity Catalog 和無伺服器運算。 增量刷新 僅支援無伺服器運算。
何時使用 REPLACE WHERE流程
請在以下情況中使用 REPLACE WHERE 流程:
- 無串流語意的增量批次處理:在不管理流式概念(如浮水印)的情況下,批量處理新資料列。
- 選擇性重處理:只重新計算與謂詞相符的列,其他列則不動。
-
超出標準物質化視圖功能的情境:
- 保留期限比來源表更長的目標表
- 當維度表變更時防止重新計算
- 結構演化而不重新計算整個歷史
建立替換 WHERE 流程
將 FLOW REPLACE WHERE 條款與 CREATE OR REFRESH STREAMING TABLE 一起內嵌使用:
CREATE OR REFRESH STREAMING TABLE orders_enriched
SCHEDULE EVERY 1 DAY
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;
在刷新過程中,目標資料表中所有與謂詞相符的列會被刪除,針對同一謂詞範圍重新計算來源查詢,並插入新的結果。 在此範例中,過去 7 天內的所有資料列都會從中刪除 orders_enriched ,並使用來源查詢重新計算。
你不需要在來源查詢中加入謂詞。 管線引擎在讀取來源時會自動套用這個值。
Note
BY NAME 是必要的。 它確保欄位是依名稱而非位置相匹配。
回填歷史資料
若要在排程刷新之外將歷史列或修正資料列寫入目標資料表,請根據歷史資料所在位置選擇兩種機制:
-
述詞覆寫:使用
REFRESH STREAMING TABLE ... WHERE重新執行該流程的來源查詢,以套用一次性的述詞範圍。 當歷史資料與增量資料來自相同來源時使用。 - DML 陳述式:直接插入到目標資料表,繞過流程。 當歷史資料與增量資料位於不同來源時使用。
述詞覆寫
在不修改資料表定義的情況下,覆寫流程的 REPLACE WHERE 謂詞一次刷新。 謂詞覆寫是一次性的,只會套用於該次重新整理,且不會影響後續排定的重新整理。
使用 WHERE 子句搭配 REFRESH STREAMING TABLE:
REFRESH STREAMING TABLE orders_enriched WHERE date BETWEEN '2020-01-01' AND '2024-12-31';
此流程只會刪除符合覆寫條件述詞的資料列,並根據來源資料重新計算這些資料列。 流程的靜態謂詞保持不變,因此下一次排程刷新仍使用原始謂詞。
你可以將覆寫設定與 SYNC 或 ASYNC 結合。 例如,若要在背景中啟動重新整理,並立即傳回管線更新的連結:
REFRESH STREAMING TABLE orders_enriched WHERE date >= date_add(current_date(), -30) ASYNC;
Note
WHERE 子句僅支援用於以 FLOW REPLACE WHERE 子句建立的串流資料表。 在其他串流表中指定該資料會回錯。 請參閱 REFRESH (MATERIALIZED VIEW 或 STREAMING TABLE)。
DML 陳述式
直接在目標資料表上執行 DML 語句,載入來自與流程不同來源的資料列:
INSERT INTO orders_enriched
SELECT *
FROM orders_enriched_legacy
WHERE date < '2025-01-01';
完整刷新行為
REPLACE WHERE 流程的完整重新整理會僅使用該流程所定義的述詞條件,重新執行來源查詢。 它會永久刪除目標中所有不符合該謂詞的列,包括先前插入超出定義謂詞範圍的謂詞覆寫或 DML 陳述的列。
Warning
完整重新整理會清除所有現有資料,並僅依該流程已定義的條件重新執行流程。 如果管線已使用 7 天的述詞條件執行了一年,完整重新整理後,資料表將只包含最近 7 天的資料。 所有舊的列都會永久刪除。
REFRESH STREAMING TABLE orders_enriched FULL;
為防止資料表被完全刷新,請將資料表屬性 pipelines.reset.allowed 設為 false:
CREATE OR REFRESH STREAMING TABLE orders_enriched
TBLPROPERTIES (pipelines.reset.allowed = 'false')
FLOW REPLACE WHERE date >= date_add(current_date(), -7) BY NAME
...
增量更新
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 |
謂詞無法下推至任何資料來源、目前視窗包含前一個謂詞尚未處理的資料列,或此次執行使用了謂詞覆寫。 |
Examples
以下範例展示了常見的替換 WHERE 流模式。
範例 1:從有限保留來源保留歷史彙總
此範例可無限期保留每日彙總資料,即使原始資料已因超過來源資料表的保留期(3 天)而被清除:
CREATE OR REFRESH 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;
範例 2:當維度表變更時防止重新計算
當維度屬性改變時,此範例保持歷史事實列不變:
CREATE OR REFRESH STREAMING TABLE fact_dim_join
FLOW REPLACE WHERE 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;
若使用者區域變更,僅會重新計算最近的資料列。 歷史行數保留當時的區域價值。
範例 3:新增指標而不重新計算完整歷史
此範例展示了如何演化表格定義,並僅回填目標範圍:
定義初始表格:
CREATE OR REFRESH 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;更新查詢以補充
uniq_users:CREATE OR REFRESH 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;早於 7 天視窗的資料列包含用於
NULL的uniq_users。
範例 4:在回填完整歷史前,先在一個小視窗上進行迭代
此範例展示了如何在處理完整歷史範圍前,先驗證一個小型資料視窗的查詢邏輯。
先以一小段時間驗證指標,並以較低的運算成本反覆調整業務邏輯:
CREATE OR REFRESH 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;
短視窗每次重新整理時只會重新計算最近 7 天的資料,因此在決定進行完整歷史重算前,可先視需要多次修改查詢。
查詢完成後,使用 DML 回填完整的歷史範圍:
INSERT INTO revenue_attribution
SELECT
event_date,
campaign_id,
SUM(revenue) AS total_revenue
FROM marketing_events
WHERE event_date < date_add(current_date(), -7)
GROUP BY ALL;