メインコンテンツまでスキップ

パイプラインによるヒストリカルデータの埋め戻し

データエンジニアリングでは、 バックフィルと は、現在のデータまたはストリーミング データを処理するために設計されたパイプ データラインを通じて履歴データを遡及的に処理するプロセスを指します。

通常、これは既存のテーブルにデータを送信する別のフローです。次の図は、パイプライン内のブロンズ テーブルにヒストリカル データを送信するバックフィル フローを示しています。

既存のワークフローに履歴データを追加するバックフィル フロー

バックフィルが必要になる可能性があるシナリオ:

  • レガシー システムからの履歴データを処理して機械学習 ( ML ) モデルをトレーニングしたり、履歴傾向分析ダッシュボードを構築したりできます。
  • アップストリーム データ ソースのデータ品質の問題のため、データのサブセットを再処理します。
  • ビジネス要件が変更されたため、最初のパイプラインでカバーされていなかった別の期間のデータをバックフィルする必要があります。
  • ビジネス ロジックが変更されたため、履歴データと現在のデータの両方を再処理する必要があります。

使用するバックフィルフローは、ターゲットテーブルとソースデータによって異なります。信頼できるスナップショットを持つ AUTO CDC slowly changing dimensions (SCD) タイプ1のターゲットの場合は、一回限りの AUTO CDC FROM SNAPSHOT フローを使用します。履歴の変更を再生する SCD マイグレーションの場合は、一回限りの AUTO CDC フローを使用します。

追加専用のバックフィルの場合: ONCE オプションを指定した専用の追加フローを使用して、追加専用のストリーミングテーブルをバックフィルします。ONCE オプションの情報については、append_flow または CREATE FLOW (パイプライン) を参照してください。

履歴データをストリーミング テーブルにバックフィルする際の考慮事項​

  • 通常、データをブロンズレイヤーのストリーミングテーブルに追加します。下流のシルバーレイヤーとゴールドレイヤーが、ブロンズレイヤーから新しいデータを取り込みます。
  • 同じデータが複数回追加された場合に、パイプラインが重複データを適切に処理できることを確認します。
  • 履歴データ スキーマが現在のデータ スキーマと互換性があることを確認してください。
  • データ量と必要な処理のサービスレベルアグリーメント(SLA)を考慮し、それに応じてクラスターとバッチサイズを設定します。

例:既存のパイプラインにバックフィルを追加する​

この例では、2025年1月1日を開始日として、クラウドストレージのソースから生のイベント登録データを取り込むパイプラインがあるとします。後になって、ダウンストリームのレポートや分析のユースケースのために、過去3年分のヒストリカルデータをバックフィルしたいことに気づきます。すべてのデータは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
)

ここでは、 modifiedAfter Auto Loader オプションを使用して、クラウド ストレージ パスからのすべてのデータを処理していないことを確認します。増分処理はその境界で中断されます。

ヒント

Kafka 、 Kinesis 、 Azure Event Hubs などの他のデータソースには、同じ動作を実現するための同等のリーダー オプションがあります。

過去3年間のデータのバックフィル​

ここで、以前のデータをバックフィルするために 1 つ以上のフローを追加します。この例では、次のステップを取り上げます。

  • append onceフローを使用します。これにより、最初のバックフィル後に実行を継続せずに、1 回限りのバックフィルが実行されます。コードはパイプライン内に残り、パイプラインが完全に更新されると、バックフィルが再実行されます。
  • 各年に 1 つずつ、合計 3 つのバックフィル フローを作成します (この場合、データはパス内で年ごとに分割されます)。Python ではフローの作成をパラメーター化しますが、SQL ではフローごとに 1 回ずつ、コードを 3 回繰り返します。

独自のプロジェクトに取り組んでいて、サーバレス コンピュートを使用していない場合は、パイプラインの最大ワーカーを更新するとよいでしょう。 最大ワーカーを増やすと、予想される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

この実装では、いくつかの重要なパターンが強調されています。

関心の分離

  • 増分処理はバックフィル操作とは独立しています。
  • 各フローには独自の構成と最適化設定があります。
  • 増分操作とバックフィル操作には明確な違いがあります。

制御された実行

  • ONCEオプションを使用すると、各バックフィルが 1 回だけ実行されるようになります。
  • バックフィル フローはパイプライン グラフ内に残りますが、完了するとアイドル状態になります。完全に更新されると自動的に使用できるようになります。
  • パイプライン定義には、バックフィル操作の明確な監査証跡があります。

処理の最適化

  • 処理を高速化するため、または処理を制御するために、大きなバックフィルを複数の小さなバックフィルに分割できます。
  • 拡張オートスケールを使用すると、現在のクラスター負荷に基づいてクラスター サイズが動的にスケーリングされます。

スキーマの展開

  • schemaEvolutionMode="addNewColumns"使用すると、スキーマの変更が適切に処理されます。
  • 履歴データと現在のデータにわたって一貫したスキーマ推論が行われます。
  • 新しいデータ内の新しい列は安全に処理されます。

AUTO CDC SCD Type 1 テーブルへのバックフィルの追加​

継続的なチェンジデータキャプチャ (CDC) フィードも受信する SCD Type 1 ターゲットに権威あるスナップショットを追加するには、1回限りの AUTO CDC FROM SNAPSHOT フローを使用します。スナップショットバージョンと CDC シーケンス列は、1つの順序付けドメインを形成します。新しい CDC イベントは古いスナップショットよりも優先され、新しいスナップショットは古い CDC イベントよりも優先されます。

要件​

バックフィルを追加する前に、フローが次の要件を満たしていることを確認してください。

  • ターゲットはSCDタイプ1を使用します。
  • ターゲットには、AUTO CDC FROM SNAPSHOTフローがちょうど1つと、一意の名前が付けられたAUTO CDCフローが1つ以上あります。
  • すべてのフローで、同じ順序で同数のキーを使用する必要があります。スナップショット フローのキー名の大文字と小文字は、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 を設定しないでください。
  • パイプラインはTriggerモードを使用します。このパターンでは継続的パイプラインはサポートされていません。

どちらのフロー タイプでも、SQLまたはPythonのパイプライン インターフェイスを使用できます。同じターゲット内でSQLフローとPythonフローを混在させることができます。

スナップショットは、そのバージョンにおけるソースの完全な状態を表す必要があります。ターゲット キーがスナップショットに存在しない場合、AUTO CDC FROM SNAPSHOT はその不在をスナップショット バージョンでの削除として扱います。新しいバージョンの CDC イベントは、キーを維持または復元します。

バックフィルの追加​

1 回限りのスナップショットのバックフィルを追加して CDC イベントの処理を継続するには、次のステップを実行します。

  1. パイプライン定義内の既存のターゲットテーブルと、実行中の AUTO CDC フローを維持します。
  2. 信頼できるスナップショットとそのバージョンを定義します。Python コールバックの場合、初回呼び出しではスナップショットとバージョンを返す必要があります。少なくとも 1 つのスナップショットが処理された後にのみ、None を返します。
  3. Pythonで AUTO CDC FROM SNAPSHOT を使用するか、SQLで ONCE を使用して、1つの once=True フローを追加します。既存のターゲットへのSQLバックフィルの場合は、WITH VERSION クエリーを含めます。WITH VERSION のないSQLスナップショットフローは、空のターゲットへの初期ロードのみをサポートします。
  4. Triggerモードのパイプライン更新をランして、バックフィルと継続的なCDCイベントを処理します。

次の例は、customers_cdcからの変更をインクリメンタルに処理する既存のPythonパイプラインから起動します。このパイプラインをすでにランし、customersターゲットにデータを設定していると仮定します。

Python
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,
)

2025年1月1日時点の customers_snapshot の状態で既存のターゲットをバックフィルするには、同じパイプライン定義に次のコードを追加します。既存のターゲット テーブルと AUTO CDC フローを維持します。

Python
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 フローは、2 つのシーケンシングタイプが完全に一致するように change_timestamp を TIMESTAMP にキャストします。この例では両方のフローに Python を使用していますが、いずれかのフローを SQL で定義し、同じターゲット内で SQL と Python のフローを混在させることもできます。空ではないターゲットに必要な WITH VERSION クエリーを含む SQL 構文については、CREATE FLOW (パイプライン) を参照してください。

スナップショットフローの commit が正常に完了すると、AUTO CDC フローが新しいイベントの処理を継続している間、後続の増分更新ではスナップショットフローがスキップされます。

重要

ターゲットを完全に更新すると、1回限りのスナップショット フローが再実行されます。完全更新を実行する前に、スナップショットを使用できる状態に保ち、意図した状態を引き続き表していることを確認してください。

この統合されたバックフィルパターンは、SCDタイプ2またはバイテンポラルターゲットをサポートしていません。

例:移行中にSCDターゲットをバックフィルする​

一般的な移行シナリオとして、レガシーシステムに何年分もの累積履歴がすでに存在しているものの、元の変更フィードが利用できなくなった slowly changing dimensions (SCD) テーブルがあります。元の変更イベントが失われているため、代わりにレガシーテーブル独自の履歴を新しい AUTO CDC ターゲットに 1 回再生し、その後で新しい CDC フィードをアタッチします。AUTO CDC と SCD タイプの詳細については、「AUTO CDC APIs : パイプラインによるチェンジデータキャプチャの簡素化」を参照してください。

このパターンは、進行中の AUTO CDC フローがターゲットとするのと同じストリーミングテーブルへの 1 回限りの AUTO CDC フローです。AUTO CDC ターゲットは AUTO CDC フローのみを受け入れるため、シードも AUTO CDC フローである必要があります。同じテーブルへの単純な INSERT INTO ONCE append フローは検証に失敗します:

  1. AUTO CDCフローが書き込む ターゲットストリーミングテーブルを作成します 。
  2. レガシー SCD テーブルをストリームとして読み取り、レガシーの有効性起動列によって順序付けられる AUTO CDC ONCE フローを使用して、 レガシー履歴を 1 回シードします 。レガシー行をご自身で整形するのではなく、変更イベントとしてリプレイします。AUTO CDC は SCD Type 2 ターゲットの __START_AT および __END_AT 履歴列を構築するため、それらの列に直接書き込まないでください。
  3. 進行中の AUTO CDC フローをアタッチし 、最新のチェンジフィードを読み取ります。AUTO CDC はキーごとの順序を解決するため、カットオーバーは各ビジネスキーに対して個別に保持する必要があります。つまり、各キーの最初のライブ変更は、同じキーの最後のシード変更の後にシーケンスされる必要があります。グローバルなレガシー最大値より単に大きいだけのシーケンス値であっても、個々のキーにとっては古い(stale)可能性があるため、そのキーの最初のライブ変更が無視されたり、誤って順序付けられたりすることがあります。

以下のコードは上記**ステップ**を用いた**ストリーミングテーブル**を作成します。

SQL
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タイプ、および順序付け列のデータ型が一致している必要があります。前述の例では、両方のフローがTimestampによって順序付けられており、単一の切り替え時刻を使用してシードされた履歴とライブフィードを分離しています。レガシーテーブルがライブフィードとは異なる型の値で順序付けられている場合は、型が一致するようにいずれかをキャストしてください。

同じシェイプがSCDタイプ1のターゲットでも機能します。両方のフローでSTORED AS SCD TYPE 2をSTORED AS SCD TYPE 1に変更すると、ターゲットはキーごとに現在の行のみを保持します。いずれかのシェイプに依存する前に、シードされたキーに対する最初のライブ変更が正確に1つの新しいバージョンを生成し、以前のバージョンを正しくクローズすることをキーのサンプルで検証してください。キーごとのシーケンスのギャップは、通常そのステップで発生します。

その他のリソース​