使用管線回填歷史資料

在資料工程中, 回填 是指透過專為處理目前或串流資料而設計的資料管道追溯處理歷史資料的過程。

通常,這是將資料傳送至現有資料表的個別流程。 下圖顯示將歷史資料傳送至管線中青銅資料表的回填流程。

回填流程:將歷史資料新增至現有工作流程

可能需要回填的一些案例:

  • 處理來自舊版系統的歷史資料,以訓練機器學習 (ML) 模型或建置歷史趨勢分析儀表板。
  • 由於上游資料來源的資料品質問題,重新處理資料子集。
  • 您的業務需求已變更,您需要回填初始管線未涵蓋的不同時段的資料。
  • 您的商務邏輯已變更,您需要重新處理歷程資料和目前資料。

你使用的回填流取決於目標資料表和來源資料。 對於具有權威快照的 AUTO CDC SCD Type 1 目標,請使用一次性 AUTO CDC FROM SNAPSHOT 流程。 對於重播歷史變更的 SCD 遷移,請使用一次性 AUTO CDC 流程。

對於僅附加回填:使用含有 ONCE 選項的專門附加流程來回填僅附加的串流表。 如需有關 選項的詳細資訊,請參閱 append_flow 或 ONCE。

將歷程記錄資料回填至串流資料表時的考量事項

  • 一般而言,請將資料附加至 Bronze 串流表格。 下游的銀層和金層則從青銅層接收到新資料。
  • 請確定您的管線可以正常處理重複資料,以防相同的資料被多次附加。
  • 請確定歷程資料結構描述與目前的資料結構描述相容。
  • 請考慮資料量大小及所需的處理服務等級協議(SLA),並相應配置叢集與批次大小。

範例:為現有管線新增資料回填

在此範例中,假設您有一條管線,可從雲端儲存來源擷取原始事件登記資料,從 2025 年 1 月 1 日開始。 你後來才發現,你想回填過去三年的歷史資料,作為後續報告和分析的應用案例。 所有資料都位於一個位置,以 JSON 格式按年、月和日分割。

初始管線

以下是從雲端儲存體累加擷取原始事件註冊資料的起始管線程式碼。

Python

from pyspark import pipelines as dp

source_root_path = spark.conf.get("registration_events_source_root_path")
begin_year = spark.conf.get("begin_year")
incremental_load_path = f"{source_root_path}/*/*/*"

# create a streaming table and the default flow to ingest streaming events
@dp.table(name="registration_events_raw", comment="Raw registration events")
def ingest():
    return (
        spark
        .readStream
        .format("cloudFiles")
        .option("cloudFiles.format", "json")
        .option("cloudFiles.inferColumnTypes", "true")
        .option("cloudFiles.maxFilesPerTrigger", 100)
        .option("cloudFiles.schemaEvolutionMode", "addNewColumns")
        .option("modifiedAfter", "2025-01-01T00:00:00.000+00:00")
        .load(incremental_load_path)
        .where(f"year(timestamp) >= {begin_year}") # safeguard to not process data before begin_year
    )

SQL

-- create a streaming table and the default flow to ingest streaming events
CREATE OR REFRESH STREAMING LIVE TABLE registration_events_raw AS
SELECT * FROM read_files(
  "/Volumes/gc/demo/apps_raw/event_registration/*/*/*",
  format => "json",
  inferColumnTypes => true,
  maxFilesPerTrigger => 100,
  schemaEvolutionMode => "addNewColumns",
  modifiedAfter => "2024-12-31T23:59:59.999+00:00"
)
WHERE year(timestamp) >= '2025'; -- safeguard to not process data before begin_year

在這裡,我們使用自動加載器選項 modifiedAfter 來確保我們不會處理來自雲存儲路徑的所有數據。 增量處理會在該界限處截斷。

小提示

其他資料來源 (例如 Kafka、Kinesis 和 Azure 事件中樞) 具有對等的讀取器選項,可達到相同的行為。

將過去 3 年的資料回填補足

現在您想要新增一或多個流程來回填先前的資料。 在此範例中,請執行下列步驟:

  • 使用 append once 流程。 這會執行一次性回填,而不會在第一次回填之後繼續執行。 程式碼會保留在您的管線中,如果管線已完全重新整理,則會重新執行回填。
  • 建立三個回填流程,每年一個 (在此情況下,資料會在路徑中依年份分割)。 對於 Python,我們參數化流程的建立,但在 SQL 中,我們重複程式碼三次,每個流程一次。

如果您正在處理自己的專案,並且未使用無伺服器運算,您可能想要更新管線的最大工作者數量。 增加工作者數目上限可確保您有資源來處理歷程記錄資料,同時繼續在預期的 SLA 內處理目前的串流資料。

小提示

如果您使用無伺服器運算搭配增強型自動擴展 (預設值),則當負載增加時,叢集的大小會自動增加。

Python

from pyspark import pipelines as dp

source_root_path = spark.conf.get("registration_events_source_root_path")
begin_year = spark.conf.get("begin_year")
backfill_years = spark.conf.get("backfill_years") # e.g. "2024,2023,2022"
incremental_load_path = f"{source_root_path}/*/*/*"

# meta programming to create append once flow for a given year (called later)
def setup_backfill_flow(year):
    backfill_path = f"{source_root_path}/year={year}/*/*"
    @dp.append_flow(
        target="registration_events_raw",
        once=True,
        name=f"flow_registration_events_raw_backfill_{year}",
        comment=f"Backfill {year} Raw registration events")
    def backfill():
        return (
            spark
            .read
            .format("json")
            .option("inferSchema", "true")
            .load(backfill_path)
        )

# create the streaming table
dp.create_streaming_table(name="registration_events_raw", comment="Raw registration events")

# append the original incremental, streaming flow
@dp.append_flow(
        target="registration_events_raw",
        name="flow_registration_events_raw_incremental",
        comment="Raw registration events")
def ingest():
    return (
        spark
        .readStream
        .format("cloudFiles")
        .option("cloudFiles.format", "json")
        .option("cloudFiles.inferColumnTypes", "true")
        .option("cloudFiles.maxFilesPerTrigger", 100)
        .option("cloudFiles.schemaEvolutionMode", "addNewColumns")
        .option("modifiedAfter", "2024-12-31T23:59:59.999+00:00")
        .load(incremental_load_path)
        .where(f"year(timestamp) >= {begin_year}")
    )

# parallelize one time multi years backfill for faster processing
# split backfill_years into array
for year in backfill_years.split(","):
    setup_backfill_flow(year) # call the previously defined append_flow for each year

SQL

-- create the streaming table
CREATE OR REFRESH STREAMING TABLE registration_events_raw;

-- append the original incremental, streaming flow
CREATE FLOW
  registration_events_raw_incremental
AS INSERT INTO
  registration_events_raw BY NAME
SELECT * FROM STREAM read_files(
  "/Volumes/gc/demo/apps_raw/event_registration/*/*/*",
  format => "json",
  inferColumnTypes => true,
  maxFilesPerTrigger => 100,
  schemaEvolutionMode => "addNewColumns",
  modifiedAfter => "2024-12-31T23:59:59.999+00:00"
)
WHERE year(timestamp) >= '2025';


-- one time backfill 2024
CREATE FLOW
  registration_events_raw_backfill_2024
AS INSERT INTO ONCE
  registration_events_raw BY NAME
SELECT * FROM read_files(
  "/Volumes/gc/demo/apps_raw/event_registration/year=2024/*/*",
  format => "json",
  inferColumnTypes => true
);

-- one time backfill 2023
CREATE FLOW
  registration_events_raw_backfill_2023
AS INSERT INTO ONCE
  registration_events_raw BY NAME
SELECT * FROM read_files(
  "/Volumes/gc/demo/apps_raw/event_registration/year=2023/*/*",
  format => "json",
  inferColumnTypes => true
);

-- one time backfill 2022
CREATE FLOW
  registration_events_raw_backfill_2022
AS INSERT INTO ONCE
  registration_events_raw BY NAME
SELECT * FROM read_files(
  "/Volumes/gc/demo/apps_raw/event_registration/year=2022/*/*",
  format => "json",
  inferColumnTypes => true
);

此實作會強調數個重要模式。

關注點分離

  • 增量處理與回填作業無關。
  • 每個流程都有自己的組態和最佳化設定。
  • 增量作業和回填作業之間有明顯的區別。

受控執行

  • 使用該 ONCE 選項可確保每個回填只執行一次。
  • 回填流程會保留在管線圖表中,但在完成後會變成閒置。 它已準備好在完全刷新時自動使用。
  • 管線定義中有回填作業的明確稽核追蹤。

處理最佳化

  • 您可以將大型回填區分割成多個較小的回填區,以加快處理速度,或控制處理。
  • 使用增強型自動調整會根據目前的叢集負載動態調整叢集大小。

綱要演進

  • 使用 schemaEvolutionMode="addNewColumns" 能夠優雅地處理結構描述變更。
  • 您在歷史和目前資料之間具有一致的結構描述推斷。
  • 在較新的資料中可以安全地處理新資料行。

為 AUTO CDC SCD 第一型資料表新增回填

使用一次性 AUTO CDC FROM SNAPSHOT 流程,將權威快照加入同時接收持續變更資料擷取(CDC)資料流的 SCD Type 1 目標。 快照版本與 CDC 定序欄形成同一個排序範圍。 較新的 CDC 事件優先於較舊的快照,而較新的快照則優先於較舊的 CDC 事件。

Requirements

在添加回填土前,請確保流程符合以下要求:

  • 目標使用 SCD Type 1。
  • 目標僅有一個 AUTO CDC FROM SNAPSHOT 流,且有一個或多個唯一命名 AUTO CDC 的流。
  • 所有流程都以相同順序使用相同數量的按鍵。 快照流程的鍵名會與 AUTO CDC 鍵名進行不區分大小寫的比較。 多個 AUTO CDC 流程必須使用相同的鍵名和大小寫。
  • 快照版本與每個 CDC 定序欄位的資料型態完全相同。
  • AUTO CDC FROM SNAPSHOT流程並未定義預期結果。
  • 這些 AUTO CDC 流不使用 IGNORE NULL UPDATES。 在 Python 中,不要設定 ignore_null_updates、 ignore_null_updates_column_list、 或 ignore_null_updates_except_column_list。
  • 管線使用觸發模式。 此模式不支援連續管線。

兩種流程類型皆可使用 SQL 或 Python 管線介面。 你可以在同一個目標中混合 SQL 和 Python 流程。

快照必須代表來源在其版本下的完整狀態。 若快照中缺少目標金鑰,AUTO CDC FROM SNAPSHOT 則將該缺失視為在該快照版本刪除。 使用較新版本的 CDC 事件會保留或還原該金鑰。

新增回填

若要新增一次性快照回填並繼續處理 CDC 事件,請使用以下步驟:

  1. 在管線定義中保留現有的目標表及其正在進行中的AUTO CDC流程。
  2. 定義權威性快照及其版本。 對於 Python 回調,第一次調用必須回傳快照和版本。 只有在至少處理一個快照後才回傳 None 。
  3. 加入一個ONCEAUTO CDC FROM SNAPSHOT流程,在 Python 或 once=True SQL 中。 若要對現有目標進行 SQL 回填,請包含 WITH VERSION 查詢。 不支援 WITH VERSION 的 SQL 快照流程只支援對空目標進行初始載入。
  4. 執行觸發的管線更新以處理回填資料和持續進行的 CDC 事件。

以下範例從一個現有的 Python 管線開始,該管線逐步處理來自 customers_cdc的變更。 假設你已經運行了這條管線,並已將資料填入 customers 目標:

from pyspark import pipelines as dp
from pyspark.sql.functions import col, expr

@dp.view
def customers_cdc():
    return (
        spark.readStream.table("main.bronze.customers_cdc")
        .withColumn("change_timestamp", col("change_timestamp").cast("timestamp"))
    )


dp.create_streaming_table("customers")

dp.create_auto_cdc_flow(
    name="customers_incremental_cdc",
    target="customers",
    source="customers_cdc",
    keys=["customer_id"],
    sequence_by=col("change_timestamp"),
    apply_as_deletes=expr("operation = 'DELETE'"),
    except_column_list=["operation", "change_timestamp"],
    stored_as_scd_type=1,
)

為了以 customers_snapshot 在 2025 年 1 月 1 日的狀態回填該現有目標,請將以下程式碼加入相同的管線定義。 保持現有目標表與 AUTO CDC 流程:

from datetime import datetime, timezone
from typing import Optional, Tuple

from pyspark.sql import DataFrame

backfill_version = datetime(2025, 1, 1, tzinfo=timezone.utc)


def backfill_snapshot_and_version(
    latest_snapshot_version: Optional[datetime],
) -> Optional[Tuple[DataFrame, datetime]]:
    if latest_snapshot_version is None:
        return (spark.read.table("main.legacy.customers_snapshot"), backfill_version)
    return None


dp.create_auto_cdc_from_snapshot_flow(
    target="customers",
    source=backfill_snapshot_and_version,
    keys=["customer_id"],
    stored_as_scd_type=1,
    once=True,
)

回調必須在第一次調用時回傳快照與版本。 如果在處理任何快照前回傳 None ,管線更新即告失敗。 快照處理完成後,回傳 None 表示沒有其他可用的快照。

快照版本是 Python datetime 型別,對應 Spark SQL TIMESTAMP 類型。 現有的 AUTO CDC 流程會將 change_timestamp 轉型為 TIMESTAMP,使兩種序列類型完全相符。 範例中兩種流程都用 Python,但你可以用 SQL 定義任一流程,並將 SQL 和 Python 流程混合在同一個目標中。 關於 SQL 語法,包括非空目標所需的 WITH VERSION 查詢,請參見 CREATE FLOW (pipelines)。

快照流程成功提交後,後續的增量更新會跳過快照,而 AUTO CDC 流程則持續處理新事件。

Important

目標的完整刷新會重新執行一次性快照流程。 在進行完整刷新前,請保持快照的可用性,並確保它仍代表預期狀態。

此統一回填圖樣不支援 SCD 第二型或雙時空目標。

範例:在遷移期間回填 SCD 目標

一個常見的遷移情境是,一個緩慢變動維度(SCD)表,該表已存在於累積多年的歷史遺留系統中,但其原始變更資料已不再存在。 因為原本的變更事件已經消失,你改為先把舊有表格的歷史重播到新 AUTO CDC 目標一次,然後再附加一個新的 CDC 資料流,後續持續接收。 若要深入了解 AUTO CDC 和 SCD 類型,請參閱 AUTO CDC API:透過管線簡化變更資料擷取。

該模式是 AUTO CDC 一次性流入與持續 AUTO CDC 流量目標相同的串流表。 AUTO CDC 目標僅接受 AUTO CDC 流,因此,種子也必須是 AUTO CDC 流。 一個單純附加到同一個資料表的 INSERT INTO ONCE 流程未通過驗證:

  1. AUTO CDC,讓您的流程寫入其中。
  2. 以一個流程一次性植入舊版歷史AUTO CDC ONCE,將舊版 SCD 資料表讀成串流,並依舊版有效起始欄位排序。 將舊有資料列重播為變更事件,而不是自行整理這些資料列。 AUTO CDC 會為 SCD Type 2 目標建立 __START_AT 和 __END_AT 歷史欄位,因此不要直接寫入這些欄位。
  3. 附加持續讀取最新變更摘要的AUTO CDC流程 AUTO CDC 會依每個鍵決定順序,因此切換作業必須對每個業務鍵分別成立:每個鍵的第一次上線變更,其順序都必須排在該鍵最後一次初始匯入的變更之後。 某個序列值即使只是晚於全域舊版最大值,就個別鍵而言仍可能已過時,而該鍵的首次有效變更因此會被忽略或排序不正確。

以下程式碼會建立一個串流表,依照上述步驟進行:

CREATE OR REFRESH STREAMING TABLE customers_history;

-- One-time seed: replay the legacy history as change events
CREATE FLOW customers_history_seed
AS AUTO CDC ONCE INTO customers_history
FROM stream(legacy.customers_scd2)
KEYS (customer_id)
SEQUENCE BY valid_from
STORED AS SCD TYPE 2;

-- Ongoing live CDC into the same target
CREATE FLOW customers_history_cdc
AS AUTO CDC INTO customers_history
FROM stream(customers_cdc_bronze)
KEYS (customer_id)
SEQUENCE BY change_timestamp
STORED AS SCD TYPE 2;

兩個流程必須在鍵、SCD 類型及排序欄位的資料型態上達成共識。 在前述範例中,兩個流程都依時間戳記排序,並使用單一切換時間來區分預先植入的歷史資料與即時資料流。 如果舊有資料表的序列值與即時資料不同,就投射其中一表,讓類型相符。

相同的模式也適用於 SCD Type 1 目標:將兩個流程中的 STORED AS SCD TYPE 2 改為 STORED AS SCD TYPE 1,目標則只保留每個鍵的目前資料列。 在依賴任一形狀之前,先驗證一組金鑰樣本,確認種子金鑰的第一次即時變更恰好產生一個新版本,且正確關閉前一個版本。 通常在那一步會出現每個鍵的序列空隙。

其他資源