以合併、插入和附加方式載入資料
在你轉換並清理資料後,流程的最後一步是將它們載入目標資料表。 Azure Databricks 提供三種核心載入策略: 附加、 覆寫與 合併。 每種策略都有不同的應用情境,選擇合適的策略取決於你的來源資料與現有紀錄的關聯。
在這個單元中,你會學習如何用 INSERT 新增資料列、用覆寫操作替換資料,以及使用 MERGE 同步變更以應付更新情境。
以 INSERT INTO 加入資料
最直接的載入操作是在現有資料表中新增資料列,且不修改當前資料。 當你載入與現有紀錄不重疊的新批次資料時,請使用 INSERT INTO,例如每日交易日誌、新事件串流或增量資料擷取。
INSERT INTO 同時適用於 SQL 和 PySpark。 以下範例說明如何為銷售資料表新增資料列。
-- Append new records using VALUES
INSERT INTO sales.transactions (transaction_id, amount, transaction_date)
VALUES
('TXN001', 150.00, '2026-01-15'),
('TXN002', 275.50, '2026-01-15');
-- Append from a staging table or query
INSERT INTO sales.transactions
SELECT * FROM staging.daily_transactions
WHERE transaction_date = current_date();
在 PySpark 中,你透過寫入帶有 append 模式的 DataFrame 來附加資料:
# Append DataFrame to existing table
df_new_transactions.write.mode("append").saveAsTable("sales.transactions")
# Or use insertInto for existing tables
df_new_transactions.write.insertInto("sales.transactions")
小提示
在資料附加時,請確保你的來源資料在載入前已完成去重處理。 插入 INTO 不會檢查重複——如果你載入同一批次兩次,就會產生重複的列。
將資料用 INSERT OVERWRITE 取代
有時候你需要替換現有資料,而不是新增資料。 INSERT OVERWRITE 會在插入新列前截斷目標資料表或分割區。 此方法適用於重新處理失敗批次、刷新查詢表或重建聚合。
-- Replace entire table contents
INSERT OVERWRITE sales.daily_summary
SELECT
region,
SUM(amount) as total_sales,
COUNT(*) as transaction_count
FROM sales.transactions
WHERE transaction_date = current_date()
GROUP BY region;
對於分割資料表,你可以覆蓋特定分割區,而保留其他分割區:
-- Replace only January 2026 partition
INSERT OVERWRITE sales.monthly_report
PARTITION (report_month = '2026-01')
SELECT region, product_category, SUM(amount) as revenue
FROM sales.transactions
WHERE transaction_date BETWEEN '2026-01-01' AND '2026-01-31'
GROUP BY region, product_category;
Delta Lake 也支援帶有REPLACE WHERE子句的選擇性覆寫,該子句會在插入前刪除符合條件的列:
-- Replace transactions for a specific date range
INSERT INTO sales.transactions
REPLACE WHERE transaction_date BETWEEN '2026-01-01' AND '2026-01-07'
SELECT * FROM staging.corrected_transactions;
在 PySpark 中,使用 overwrite 完整資料表替換模式:
# Overwrite entire table
df_refreshed.write.mode("overwrite").saveAsTable("sales.daily_summary")
# Overwrite specific partition
df_january.write.mode("overwrite").partitionBy("report_month").saveAsTable("sales.monthly_report")
這很重要
SQL INSERT OVERWRITE PARTITION 語法與 partitionOverwriteMode=dynamic 結合使用時,僅限於傳統計算- 無法在 Databricks SQL 倉儲或無伺服器計算上運作。 使用 partitionOverwriteMode=dynamic 進行的 Python 和 Scala DataFrame 寫入,可在所有計算類型上運作。 對於新的管線,Databricks 建議使用 REPLACE USING 而非分割區覆寫,因為分割區覆寫在分區變更時可能會使用過時的資料。
動態地使用 REPLACE USING 替換資料
REPLACE USING 是動態資料覆寫的推薦方法。 它適用於所有運算類型——包括 Databricks、SQL 倉庫和無伺服器運算——且不需要 Spark 會話設定。 與 partitionOverwriteMode不同, REPLACE USING 支援分割表、非分割表及具液態叢集的資料表。
-- Replace only the partitions touched by the incoming data
INSERT INTO sales.monthly_report
REPLACE USING (report_month)
SELECT region, product_category, SUM(amount) AS revenue, report_month
FROM staging.corrected_transactions
GROUP BY region, product_category, report_month;
REPLACE USING 子句會自動刪除sales.monthly_report中的資料列,在report_month 上符合傳入資料列,然後插入新資料列。 不符合指定欄位的列則保持不變。 在 Databricks Runtime 16.3–17.1 中,欄位 USING 必須是資料表完整的分割欄位集合(舊有行為)。 在 Databricks Runtime 17.2 及以上版本中,支援任何資料表類型的欄位——包括未分割的資料表及具液態叢集的資料表。
當你需要更靈活的匹配方式——例如根據自訂條件替換列或將 NULL 值視為相等——則改用 REPLACE ON :
-- Replace rows matching a user-defined condition across source and target
INSERT INTO sales.transactions
REPLACE ON (target.transaction_id = source.transaction_id)
SELECT * FROM staging.corrected_transactions AS source;
REPLACE ON 需要 Databricks 執行環境 17.1 及以上版本。 當你需要基於條件的替換且沒有來源表參考時,請使用 REPLACE WHERE (前面提到的)。
合併資料以進行更新或插入操作
現實世界的資料管線通常需要在插入新紀錄的同時更新現有紀錄。 MERGE 陳述式處理這種模式——通常稱為「upsert」——透過將來源列與目標列匹配,並根據是否存在匹配來執行不同動作。
MERGE 對於緩慢改變維度、同步來源系統資料,以及處理變更資料擷取(CDC)資料流至關重要。
MERGE INTO customers AS target
USING customer_updates AS source
ON target.customer_id = source.customer_id
WHEN MATCHED THEN
UPDATE SET
target.email = source.email,
target.phone = source.phone,
target.last_updated = current_timestamp()
WHEN NOT MATCHED THEN
INSERT (customer_id, name, email, phone, created_date, last_updated)
VALUES (source.customer_id, source.name, source.email, source.phone,
current_timestamp(), current_timestamp());
此查詢會更新現有客戶的電子郵件與電話,同時插入全新的客戶紀錄。 ON 子句定義了匹配金鑰——通常是主金鑰或商業識別碼。
與條件邏輯合併
在 WHEN 子句中加入條件,以處理更複雜的情境。 例如,你可能只更新實際變更的紀錄:
MERGE INTO products AS target
USING product_feed AS source
ON target.sku = source.sku
WHEN MATCHED AND source.price <> target.price THEN
UPDATE SET target.price = source.price, target.updated_at = current_timestamp()
WHEN MATCHED AND source.discontinued = true THEN
DELETE
WHEN NOT MATCHED THEN
INSERT *;
語法 INSERT * 會將來源中所有符合目標結構的欄位插入。
與 PySpark 合併
程式控制方面,請使用 Delta Lake Python API:
from delta.tables import DeltaTable
target_table = DeltaTable.forName(spark, "customers")
target_table.alias("target").merge(
source_df.alias("source"),
"target.customer_id = source.customer_id"
).whenMatchedUpdate(set={
"email": "source.email",
"phone": "source.phone",
"last_updated": "current_timestamp()"
}).whenNotMatchedInsert(values={
"customer_id": "source.customer_id",
"name": "source.name",
"email": "source.email",
"phone": "source.phone",
"created_date": "current_timestamp()",
"last_updated": "current_timestamp()"
}).execute()
對於來源與目標結構相符的簡單情況,請使用以下便利方法:
target_table.alias("target").merge(
source_df.alias("source"),
"target.customer_id = source.customer_id"
).whenMatchedUpdateAll(
).whenNotMatchedInsertAll(
).execute()
選擇合適的裝填策略
你選擇的載入操作取決於來源資料與目標資料之間的關係:
| Scenario | 作業 | 使用時機 |
|---|---|---|
| 僅限新紀錄 | 插入/追加 | 每日日誌、事件串流、增量擷取 |
| 完整重新整理 | 插入覆寫 | 查詢表、聚合、批次重處理失敗 |
| 混合更新與插入 | MERGE | CDC 資料流、維度表、資料同步 |
| 依謂語選擇性替換 | 替換,其中 | 修正特定日期範圍或邏輯條件 |
| 動態資料覆寫(所有運算類型) | 替換使用 | 建議用於任何資料表類型(分區、非分區、液態叢集)的列層覆寫 |
| 客製化條件列替換 | 替換開啟 | 具有 NULL 安全或複雜比對條件的資料列層級覆寫 |
這很重要
MERGE 操作需要來源列與目標列之間唯一的匹配。 若多個來源列匹配同一目標列,操作即失敗。 合併前先移除來源資料中的重複項。
了解這些載入模式能幫助你建立可靠的資料管線,維持資料完整性,同時有效處理新紀錄與變更紀錄。