レイクフローパイプラインにおける次元モデリング

次元モデリングは、ゴールドレイヤーデータをファクトテーブルと次元テーブルに整理し、アナリストやビジネスインテリジェンス(BI)ツールが効率的にクエリできるようにするための手法です。 このページでは、Lakeflowパイプラインでそのモデルを構築する方法を説明しています。

Overview

次元モデリングはデータを2種類のテーブルに分けます:

  • ファクトテーブルには 、注文、クリック、売上など、あなたが関心を持つイベントや測定値が記録されています。 各行はそのイベントの1回の発生を示し、主にキーと数値の指標で表されます。
  • ディメンションテーブルは 、顧客、商品、日付などのイベントに関する記述的な文脈を保持します。 各行は1つの法人を表します。

スタースキーマとは、中央に1つのファクトテーブルを置き、それをキーで複数の次元テーブルに連結した形のことです。 レイアウトはアナリストやBIツールがクエリしやすく、エンジニアも論理的に考えやすいです。なぜなら、各テーブルには明確な責任があるからです。

Lakeflow のパイプラインにおいて、スター スキーマはメダリオン アーキテクチャのゴールド レイヤーに自然に適合します。 ブロンズとシルバーのデータセットは取り込みとクレンジングを担当し、ゴールドは事実表やディメンション表を具現化して、下流の消費者が直接クエリできるようにします。 パイプラインがこれらのテーブルを段階的に更新してくれるため、BIレイヤーで別途抽出・変換・ロード(ETL)ステップを経ずに、スタースキーマのクエリのシンプルさを実現できます。

どのように機能するのか

パイプライン内で次元や事実をデータセットとして構築し、それぞれの変化に合わせてデータセットタイプを選択します。 ほとんどの金層モデルについて:

  • 次元テーブルはマテリアライズされたビューとして構築する(または履歴が必要な場合は、ゆっくり変化する次元(SCD)タイプ2のストリーミングテーブルとして作成してください。 マテリアライズドビューは、入力の変更に応じてクレンジング済みのシルバーデータから効率的に再計算され、ビジネスエンティティごとに1行が得られます。
  • ファクトテーブルは、シルバー層から増分で取り込むストリーミングテーブルとして構築し、ゴールドレイヤーの集計がほぼリアルタイムを維持できるようにします。 ファクトは、記述属性を重複して持つのではなく、キーによってディメンションを参照します。

2つのデータセットタイプの詳細については、「 マテリアル化されたビュー 」と 「ストリーミングテーブル」をご覧ください。 次元ごとに履歴を追跡するには、 The AUTO CDC API: Streamify Change Data capture with pipelines(パイプラインによる変更データキャプチャの簡素化)をご覧ください。

鍵と代理鍵

ソースのナチュラル キー (ソースデータに既に存在する識別子、例えば注文番号など)を好んで選びます。これはソースのナチュラルキーが安定して利用可能で、クラスタリングや結合がうまくできるからです。 ソースがIDを再利用または変更した場合にのみ、 サロゲートキー (パイプライン生成の代替識別子)に利用してください。

サロゲートキーが必要な場合は、sha2(natural_key)のようなハッシュサロゲートは避けてください。 ハッシュは意図的にランダムで、物理的に隣接する行がファイルに散らばってしまうため、液体クラスタリングやZオーダー性能に悪影響を及ぼします。 代わりに、安定なナチュラルキーから決定論的に順序保持の代理を導き出し、同じビジネスエンティティが常に同じ代理にマッピングされるようにします。 決定論キーは、次元の完全な刷新や再構築を経ても存続し、既存の事実から次元の結合がそのまま維持されます。

あるいは、上流テーブルが追加のみでフルリフレッシュされない場合に IDENTITY 列を使うこともできます。 値 IDENTITY 行挿入時に割り当てられるため、リビルドでは同じエンティティに異なるIDを再割り当てし、以前の値を保持していた事実から次元への結合を静かに切断できます。

日付ディメンション

dim_dateを、sequence()explode()で生成するシンプルな具体化されたビューとして構築し、ソースから取り込むのではなく、日付範囲にわたって生成してください。 これは静的な参照データで、計算コストも安く、モデル内の日付に基づく結合やウィンドウ作成も簡素化されます。

Examples

以下の例は、顧客次元と注文ファクトテーブルを備えた小さなスタースキーマを構築します。

次元テーブル

ディメンションテーブルは通常、クリーン化されたシルバーデータから構築された物質化されたビューで、事業体ごとに1行ずつ構成されます。以下のコードのように:

Python

from pyspark import pipelines as dp

@dp.materialized_view(name="dim_customer", comment="Customer dimension")
def dim_customer():
    return (
        spark.read.table("customers_silver")
        .select("customer_id", "customer_name", "region", "signup_date")
    )

SQL

CREATE OR REFRESH MATERIALIZED VIEW dim_customer
COMMENT "Customer dimension"
AS SELECT customer_id, customer_name, region, signup_date
FROM customers_silver;

ファクト テーブル

ファクトテーブルは測定可能な出来事を保持し、記述的属性を重複するのではなくキーで次元を参照します。 事実は主にキーと数値指標で絞り、クエリ時に記述的な詳細を引き出すためにジョインを使ってください。以下のコードのように:

Python

from pyspark import pipelines as dp

@dp.table(name="fact_orders", comment="One row per order line, keyed to dimensions")
def fact_orders():
    return (
        spark.readStream.table("orders_silver")
        .select(
            "order_id",
            "customer_id",       # foreign key to dim_customer
            "product_id",        # foreign key to dim_product
            "order_date",        # foreign key to dim_date
            "quantity",
            "amount",
        )
    )

SQL

CREATE OR REFRESH STREAMING TABLE fact_orders
COMMENT "One row per order line, keyed to dimensions"
AS SELECT
  order_id,
  customer_id,   -- foreign key to dim_customer
  product_id,    -- foreign key to dim_product
  order_date,    -- foreign key to dim_date
  quantity,
  amount
FROM STREAM(orders_silver);

ベスト プラクティス

スタースキーマが成長する中で健全に保ついくつかの実践例があります:

  • 事実はストリーミングテーブルとして、次元は具体化されたビューとして保持し、特に履歴を変更する必要がある場合は、STORED AS SCD TYPE 2AUTO CDCを組み合わせて使うのが良いです。 「AUTO CDC API: パイプラインを使用して変更データ キャプチャを簡略化する」を参照してください。
  • 下流のBIツールを使って、ゴールドの物質化されたビューを直接照会しましょう。 Lakeflowパイプラインは段階的にリフレッシュされるため、個別の報告ETLステップなしでほぼリアルタイムの結果が得られます。
  • 次元と事実を別々のフローとして同一のゴールドレイヤーにまとめることで、各データセットを一つの一貫したDAGの一部としてスケジュール、チェックポイント、更新が可能です。 「Lakeflow パイプライン フローを使用してデータを増分的に読み込んで処理する」を参照してください。

その他のリソース