1
1

Delete article

Deleted articles cannot be recovered.

Draft of this article would be also deleted.

Are you sure you want to delete this article?

BigQuery SQL・BigQuery DataFrames・Sparkで実装するBronze to Silver

1
Posted at

はじめに

メダリオンアーキテクチャでは、ソースから取り込んだデータを Bronze レイヤーに保持し、型変換や重複排除などの加工を行ったうえで、Silver レイヤーへ反映します。

本記事では、BigQuery の変更履歴を取得できる CHANGES 関数と MERGE 文を組み合わせ、Bronze テーブルの変更内容を Silver テーブルへ差分反映する処理を検証します。

実装方法として、次の3パターンを取り上げます。

  • BigQuery SQL
  • BigQuery DataFrames
  • Spark in BigQuery

検証コードは BigQuery ノートブック上で実行します。いずれの方法でも、指定期間の変更履歴から主キーごとの最新レコードを取得し、Silverテーブルのスキーマに合わせて型変換したうえで、既存レコードの更新または新規レコードの追加を行います。

CHANGES 関数は1回の呼び出しで取得できる変更履歴の期間が最大1日という制約があります。そのため、本記事では複数日分の変更履歴を取得する場合、対象期間を1日単位に分割して CHANGES 関数を実行し、それぞれの結果をまとめて処理します。

なお、Bronze テーブルが追記のみで運用される場合は、CHANGES 関数ではなく APPENDS 関数を利用する方法も考えられます。

また、取り込み日時には CHANGES 関数が返す _CHANGE_TIMESTAMP を使用します。取り込み時間パーティションの値の利用も検証しましたが、分や秒が丸められる場合があるため注意が必要です。

本記事の処理内容は過去に講師を担当した下記のコードを参考にしています。

サンプルテーブル

まず、今回の検証で使用する Bronze テーブルと Silver テーブルを作成します。

Bronze テーブルは値を文字列として保持します。一方、Silver テーブルでは各列を用途に応じたデータ型へ変換し、主キーと取り込み日時・更新日時の監査列を定義します。

変更履歴を CHANGES 関数から取得できるように、Bronze テーブルでは変更履歴を有効化します。

Bronze 作成
%%bigquery

CREATE OR REPLACE TABLE `medallion_demo.product2__bronze` (
  `Id` STRING,
  `Name` STRING,
  `ProductCode` STRING,
  `Description` STRING,
  `IsActive` STRING,
  `CreatedDate` STRING,
  `CreatedById` STRING,
  `LastModifiedDate` STRING,
  `LastModifiedById` STRING,
  `SystemModstamp` STRING,
  `Family` STRING,
  `ExternalDataSourceId` STRING,
  `ExternalId` STRING,
  `DisplayUrl` STRING,
  `QuantityUnitOfMeasure` STRING,
  `IsDeleted` STRING,
  `IsArchived` STRING,
  `LastViewedDate` STRING,
  `LastReferencedDate` STRING,
  `StockKeepingUnit` STRING
);

ALTER TABLE `medallion_demo.product2__bronze`
SET OPTIONS (
  enable_change_history = TRUE
);
Bronze 初期データ
%%bigquery

TRUNCATE TABLE `medallion_demo.product2__bronze`;

INSERT INTO `medallion_demo.product2__bronze` (
  `Id`,
  `Name`,
  `ProductCode`,
  `Description`,
  `IsActive`,
  `CreatedDate`,
  `CreatedById`,
  `LastModifiedDate`,
  `LastModifiedById`,
  `SystemModstamp`,
  `Family`,
  `ExternalDataSourceId`,
  `ExternalId`,
  `DisplayUrl`,
  `QuantityUnitOfMeasure`,
  `IsDeleted`,
  `IsArchived`,
  `LastViewedDate`,
  `LastReferencedDate`,
  `StockKeepingUnit`
)
VALUES
  (
    '01t000000000001AAA',
    'ノートパソコン Standard',
    'PC-STD-001',
    '一般業務向けの標準ノートパソコン',
    'true',
    '2026-07-01T09:00:00.000+0000',
    '005000000000001AAA',
    '2026-07-20T10:30:00.000+0000',
    '005000000000002AAA',
    '2026-07-20T10:30:00.000+0000',
    'Hardware',
    NULL,
    'EXT-PRODUCT-001',
    'https://example.com/products/PC-STD-001',
    '',
    'false',
    'false',
    '2026-07-29T01:00:00.000+0000',
    '2026-07-29T01:05:00.000+0000',
    'SKU-PC-STD-001'
  ),
  (
    '01t000000000002AAA',
    'ノートパソコン Premium',
    'PC-PRM-001',
    '高性能プロセッサーを搭載した開発者向けノートパソコン',
    'true',
    '2026-07-02T09:00:00.000+0000',
    '005000000000001AAA',
    '2026-07-21T11:00:00.000+0000',
    '005000000000002AAA',
    '2026-07-21T11:00:00.000+0000',
    'Hardware',
    NULL,
    'EXT-PRODUCT-002',
    'https://example.com/products/PC-PRM-001',
    '',
    'false',
    'false',
    '2026-07-29T02:00:00.000+0000',
    '2026-07-29T02:10:00.000+0000',
    'SKU-PC-PRM-001'
  ),
  (
    '01t000000000003AAA',
    'クラウドサービス Basic',
    'SVC-BSC-001',
    '月額契約のクラウドサービス基本プラン',
    'false',
    '2026-07-03T09:00:00.000+0000',
    '005000000000001AAA',
    '2026-07-25T15:00:00.000+0000',
    '005000000000003AAA',
    '2026-07-25T15:00:00.000+0000',
    'Service',
    NULL,
    'EXT-PRODUCT-003',
    'https://example.com/products/SVC-BSC-001',
    'ライセンス',
    'false',
    'true',
    NULL,
    NULL,
    'SKU-SVC-BSC-001'
  )
;

SELECT
  *
FROM medallion_demo.product2__bronze;
Bronze 追加データ
%%bigquery

INSERT INTO `medallion_demo.product2__bronze` (
  `Id`,
  `Name`,
  `ProductCode`,
  `Description`,
  `IsActive`,
  `CreatedDate`,
  `CreatedById`,
  `LastModifiedDate`,
  `LastModifiedById`,
  `SystemModstamp`,
  `Family`,
  `ExternalDataSourceId`,
  `ExternalId`,
  `DisplayUrl`,
  `QuantityUnitOfMeasure`,
  `IsDeleted`,
  `IsArchived`,
  `LastViewedDate`,
  `LastReferencedDate`,
  `StockKeepingUnit`
)
VALUES (
    '01t000000000001AAA',
    'ノートパソコン Standard',
    'PC-STD-001',
    '一般業務向けの標準ノートパソコン',
    'true',
    '2026-07-01T09:00:00.000+0000',
    '005000000000001AAA',
    '2026-07-20T10:30:00.000+0000',
    '005000000000002AAA',
    '2026-07-20T10:30:00.000+0000',
    'Hardware',
    NULL,
    'EXT-PRODUCT-001',
    'https://example.com/products/PC-STD-001',
    '',
    'false',
    'false',
    '2026-07-30T01:00:00.000+0000',
    '2026-07-30T01:05:00.000+0000',
    'SKU-PC-STD-001'
  ),
  (
    '01t000000000004AAA',
    'クラウドサービス Premium',
    'SVC-PRM-001',
    '大規模利用向けのクラウドサービス上位プラン',
    'true',
    '2026-07-04T09:00:00.000+0000',
    '005000000000001AAA',
    '2026-07-30T06:00:00.000+0000',
    '005000000000003AAA',
    '2026-07-30T06:00:00.000+0000',
    'Service',
    NULL,
    'EXT-PRODUCT-004',
    'https://example.com/products/SVC-PRM-001',
    'ライセンス',
    'false',
    'false',
    '2026-07-30T06:01:00.000+0000',
    '2026-07-30T06:02:00.000+0000',
    'SKU-SVC-PRM-001'
  );

SELECT
  *
FROM medallion_demo.product2__bronze;
Silver 作成
%%bigquery

CREATE OR REPLACE TABLE `medallion_demo.product2__silver` (
  `Id` STRING NOT NULL,
  `Name` STRING,
  `ProductCode` STRING,
  `Description` STRING,
  `IsActive` BOOL,
  `CreatedDate` TIMESTAMP,
  `CreatedById` STRING,
  `LastModifiedDate` TIMESTAMP,
  `LastModifiedById` STRING,
  `SystemModstamp` TIMESTAMP,
  `Family` STRING,
  `ExternalDataSourceId` STRING,
  `ExternalId` STRING,
  `DisplayUrl` STRING,
  `QuantityUnitOfMeasure` STRING,
  `IsDeleted` BOOL,
  `IsArchived` BOOL,
  `LastViewedDate` TIMESTAMP,
  `LastReferencedDate` TIMESTAMP,
  `StockKeepingUnit` STRING,
  `_ingest_timestamp` TIMESTAMP,
  `_update_timestamp` TIMESTAMP,

  PRIMARY KEY (`Id`) NOT ENFORCED
)
;

ALTER TABLE `medallion_demo.product2__silver`
SET OPTIONS (
  enable_change_history = TRUE
);

BigQuery SQL

最初に、BigQuery SQL を使用して Bronze テーブルから Silver テーブルへ差分反映する方法を確認します。

指定期間を1日単位に分割して CHANGES 関数を実行し、取得した変更履歴から主キーごとの最新レコードを抽出します。その後、SAFE_CAST で Silver テーブルのデータ型に変換し、MERGE 文を使って更新または追加を行います。

SQL だけで処理が完結するため、処理内容を BigQuery 上で直接確認しやすい実装です。

BigQuery SQL による差分連携
%%bigquery

DECLARE end_ts TIMESTAMP DEFAULT CURRENT_TIMESTAMP();
DECLARE start_ts TIMESTAMP DEFAULT TIMESTAMP_SUB(end_ts, INTERVAL 7 DAY);

MERGE `medallion_demo.product2__silver` AS slv
USING (
  WITH brz_changes AS (
    SELECT *
    FROM CHANGES(
      TABLE `medallion_demo.product2__bronze`,
      start_ts,
      TIMESTAMP_ADD(start_ts, INTERVAL 1 DAY)
    )
    UNION ALL
    SELECT *
    FROM CHANGES(
      TABLE `medallion_demo.product2__bronze`,
      TIMESTAMP_ADD(start_ts, INTERVAL 1 DAY),
      TIMESTAMP_ADD(start_ts, INTERVAL 2 DAY)
    )
    UNION ALL
    SELECT *
    FROM CHANGES(
      TABLE `medallion_demo.product2__bronze`,
      TIMESTAMP_ADD(start_ts, INTERVAL 2 DAY),
      TIMESTAMP_ADD(start_ts, INTERVAL 3 DAY)
    )
    UNION ALL
    SELECT *
    FROM CHANGES(
      TABLE `medallion_demo.product2__bronze`,
      TIMESTAMP_ADD(start_ts, INTERVAL 3 DAY),
      TIMESTAMP_ADD(start_ts, INTERVAL 4 DAY)
    )
    UNION ALL
    SELECT *
    FROM CHANGES(
      TABLE `medallion_demo.product2__bronze`,
      TIMESTAMP_ADD(start_ts, INTERVAL 4 DAY),
      TIMESTAMP_ADD(start_ts, INTERVAL 5 DAY)
    )
    UNION ALL
    SELECT *
    FROM CHANGES(
      TABLE `medallion_demo.product2__bronze`,
      TIMESTAMP_ADD(start_ts, INTERVAL 5 DAY),
      TIMESTAMP_ADD(start_ts, INTERVAL 6 DAY)
    )
    UNION ALL
    SELECT *
    FROM CHANGES(
      TABLE `medallion_demo.product2__bronze`,
      TIMESTAMP_ADD(start_ts, INTERVAL 6 DAY),
      end_ts
    )
  ),
  slv_records AS (
    SELECT
        Id,
        MAX(`_CHANGE_TIMESTAMP`) AS max_ingest_timestamp
    FROM brz_changes
    GROUP BY
        id
  )
  SELECT
    brz.`Id`,
    brz.`Name`,
    brz.`ProductCode`,
    brz.`Description`,
    SAFE_CAST(brz.`IsActive` AS BOOL) AS `IsActive`,
    SAFE_CAST(brz.`CreatedDate` AS TIMESTAMP) AS `CreatedDate`,
    brz.`CreatedById`,
    SAFE_CAST(brz.`LastModifiedDate` AS TIMESTAMP) AS `LastModifiedDate`,
    brz.`LastModifiedById`,
    SAFE_CAST(brz.`SystemModstamp` AS TIMESTAMP) AS `SystemModstamp`,
    brz.`Family`,
    brz.`ExternalDataSourceId`,
    brz.`ExternalId`,
    brz.`DisplayUrl`,
    brz.`QuantityUnitOfMeasure`,
    SAFE_CAST(brz.`IsDeleted` AS BOOL) AS `IsDeleted`,
    SAFE_CAST(brz.`IsArchived` AS BOOL) AS `IsArchived`,
    SAFE_CAST(brz.`LastViewedDate` AS TIMESTAMP) AS `LastViewedDate`,
    SAFE_CAST(brz.`LastReferencedDate` AS TIMESTAMP) AS `LastReferencedDate`,
    brz.`StockKeepingUnit`,
    brz.`_CHANGE_TIMESTAMP` AS `_ingest_timestamp`,
  FROM brz_changes AS brz
  INNER JOIN slv_records AS slvr
    ON brz.id =  slvr.id
    AND brz.`_CHANGE_TIMESTAMP` =  slvr.max_ingest_timestamp
  WHERE brz.`_CHANGE_TYPE` IN ('INSERT', 'UPDATE')
  QUALIFY ROW_NUMBER() OVER (
    PARTITION BY brz.`Id`
    ORDER BY
      brz.`_CHANGE_TIMESTAMP` DESC
  ) = 1
) AS brz
ON slv.`Id` = brz.`Id`
WHEN MATCHED AND slv.`_ingest_timestamp` < brz.`_ingest_timestamp` THEN
  UPDATE SET
    slv.`Name` = brz.`Name`,
    slv.`ProductCode` = brz.`ProductCode`,
    slv.`Description` = brz.`Description`,
    slv.`IsActive` = brz.`IsActive`,
    slv.`CreatedDate` = brz.`CreatedDate`,
    slv.`CreatedById` = brz.`CreatedById`,
    slv.`LastModifiedDate` = brz.`LastModifiedDate`,
    slv.`LastModifiedById` = brz.`LastModifiedById`,
    slv.`SystemModstamp` = brz.`SystemModstamp`,
    slv.`Family` = brz.`Family`,
    slv.`ExternalDataSourceId` = brz.`ExternalDataSourceId`,
    slv.`ExternalId` = brz.`ExternalId`,
    slv.`DisplayUrl` = brz.`DisplayUrl`,
    slv.`QuantityUnitOfMeasure` = brz.`QuantityUnitOfMeasure`,
    slv.`IsDeleted` = brz.`IsDeleted`,
    slv.`IsArchived` = brz.`IsArchived`,
    slv.`LastViewedDate` = brz.`LastViewedDate`,
    slv.`LastReferencedDate` = brz.`LastReferencedDate`,
    slv.`StockKeepingUnit` = brz.`StockKeepingUnit`,
    slv.`_ingest_timestamp` = brz.`_ingest_timestamp`,
    slv.`_update_timestamp` = current_timestamp()
WHEN NOT MATCHED THEN
  INSERT (`Id`, `Name`, `ProductCode`, `Description`, `IsActive`, `CreatedDate`, `CreatedById`, `LastModifiedDate`, `LastModifiedById`, `SystemModstamp`, `Family`, `ExternalDataSourceId`, `ExternalId`, `DisplayUrl`, `QuantityUnitOfMeasure`, `IsDeleted`, `IsArchived`, `LastViewedDate`, `LastReferencedDate`, `StockKeepingUnit`, `_ingest_timestamp`, `_update_timestamp`)
    VALUES (brz.`Id`, brz.`Name`, brz.`ProductCode`, brz.`Description`, brz.`IsActive`, brz.`CreatedDate`, brz.`CreatedById`, brz.`LastModifiedDate`, brz.`LastModifiedById`, brz.`SystemModstamp`, brz.`Family`, brz.`ExternalDataSourceId`, brz.`ExternalId`, brz.`DisplayUrl`, brz.`QuantityUnitOfMeasure`, brz.`IsDeleted`, brz.`IsArchived`, brz.`LastViewedDate`, brz.`LastReferencedDate`, brz.`StockKeepingUnit`, brz.`_ingest_timestamp`, current_timestamp())
;

BigQuery Dataframe

続いて、BigQuery DataFrames を使って同じ処理を実装します。

BigQuery DataFrames では、変更履歴の取得や最新レコードの判定、型変換といった処理を DataFrame 操作として記述します。最後に、一時テーブルへ処理結果を書き込み、BigQuery の MERGE 文を実行して Silver テーブルへ反映します。

SQL をすべて直接記述する方法と比べて、処理を関数単位に分割しやすく、複数のテーブルに適用できる共通処理として整理しやすい点が特徴です。

処理の内容

処理は、主に次の流れで構成します。

  1. Silver テーブルに定義された主キーを取得する
  2. Bronze テーブルと Silver テーブルのスキーマを取得する
  3. Bronze テーブルの変更履歴から主キーごとの最新レコードを取得する
  4. Silver テーブルのスキーマに合わせて型変換する
  5. 処理結果を一時テーブルへ書き込む
  6. 一時テーブルから Silver テーブルへ MERGE する
  7. 処理後に一時テーブルを削除する

以降では、それぞれの処理を関数に分けて実装しています。

設定
from __future__ import annotations

from collections.abc import Iterable
from datetime import datetime, timedelta, timezone
from zoneinfo import ZoneInfo

import bigframes.bigquery as bbq
import bigframes.pandas as bpd
from google.cloud import bigquery

PROJECT_ID = "<PROJECT_ID>"
LOCATION = "asia-northeast1"

bpd.options.bigquery.project = PROJECT_ID
bpd.options.bigquery.location = LOCATION

bq_client = bigquery.Client(
    project=PROJECT_ID,
    location=LOCATION,
)
主キー取得
def get_primary_keys(
    client: bigquery.Client,
    table_full_name: str,
  ) -> list[str]:
    """
    BigQueryテーブルに定義された主キー列を取得する。

    Args:
        client: BigQueryクライアント
        table_full_name: project.dataset.table 形式のテーブル名。

    Returns:
        主キー列名のリスト。
    """

    table = client.get_table(table_full_name)
    constraints = table.table_constraints

    if constraints is None or constraints.primary_key is None:
        return []

    return list(constraints.primary_key.columns)
スキーマ取得
def get_table_schema(
    client: bigquery.Client,
    table_full_name: str
  ) -> dict[str, str]:
    """
    BigQueryテーブルの列名とデータ型を取得する。

    Args:
        client: BigQueryクライアント
        table_full_name: project.dataset.table 形式のテーブル名。

    Returns:
        列名をキー、BigQueryデータ型を値とする辞書。
    """

    table = client.get_table(table_full_name)
    
    return {
        field.name: field.field_type
        for field in table.schema
        }
主キーごとの最新レコードを取得
def latest_by_keys(
    client: bigquery.Client,
    source: str,
    target_keys: Iterable[str],
    audit_cols: dict[str, str],
    end_ts: datetime | None = None,
    start_ts: datetime | None = None,
    tz: str | None = None,
    days_int: int | None = None,
  ) -> bpd.DataFrame:
    """
    指定期間の変更履歴から、主キーごとの最新レコードを取得する。

    CHANGES関数が返すINSERTおよびUPDATEのレコードを対象とし、
    主キーごとに_CHANGE_TIMESTAMPが最も新しいレコードを残す。

    CHANGES関数の取得可能期間は1回につき最大1日であるため、
    指定期間を1日単位に分割して取得する。

    取得した_CHANGE_TIMESTAMPは、Silverテーブルで使用する
    取り込み日時列に設定する。

    Args:
        client: 
            BigQueryクライアント
        source:
            変更履歴を取得するBigQueryテーブル。
            project.dataset.table 形式で指定する。
        target_keys:
            最新レコードを判定する主キー列。
        audit_cols:
            監査列名の定義。
            ingest_timestamp_colキーを含む必要がある。
        end_ts:
            変更履歴の取得終了日時。指定時刻は含まれない。
            未指定の場合は現在日時。
        start_ts:
            変更履歴の取得開始日時。
            未指定の場合はend_tsからdays_int日前。
        tz:
            日時を解釈するタイムゾーン。
            未指定の場合はUTC。
        days_int:
            start_tsを省略した場合の取得日数。
            未指定の場合は7日。

    Returns:
        主キーごとの最新レコードを保持するBigQuery DataFrame。

    Raises:
        ValueError: 主キーが指定されていない場合。
    """

    target_keys = list(target_keys)

    if not target_keys:
        raise ValueError("keysには1つ以上の主キー列を指定してください。")

    ingest_timestamp_col = audit_cols.get("ingest_timestamp_col")

    days_int = 7 if days_int is None else days_int
    tz = "UTC" if tz is None else tz
    if end_ts is None:
        end_ts = datetime.now(ZoneInfo(tz))
    if start_ts is None:
        start_ts = end_ts - timedelta(days=days_int)

    source_columns = list(
        get_table_schema(
            client=client,
            table_full_name=source,)
    )

    # CHANGES関数の予約列には、通常の列名として使用できる別名を付ける。
    change_type_col = "change_type_internal"
    change_timestamp_col = "change_timestamp_internal"
    source_columns_sql = ",\n              ".join(
        f"`{column}`"
        for column in source_columns
    )

    change_queries: list[str] = []
    range_start = start_ts

    while range_start < end_ts:
        range_end = min(
            range_start + timedelta(days=1),
            end_ts,
        )
        range_start_literal = range_start.isoformat(
            timespec="microseconds"
        )
        range_end_literal = range_end.isoformat(
            timespec="microseconds"
        )
        change_queries.append(
            f"""
            SELECT
              {source_columns_sql},
              _CHANGE_TYPE AS `{change_type_col}`,
              _CHANGE_TIMESTAMP AS `{change_timestamp_col}`
            FROM CHANGES(
              TABLE `{source}`,
              TIMESTAMP('{range_start_literal}'),
              TIMESTAMP('{range_end_literal}')
            )
            """
        )
        range_start = range_end

    changes_query = "\nUNION ALL\n".join(change_queries)

    changes_df = bpd.read_gbq_query(changes_query)

    # 更新後のレコードだけを対象にする
    changes_df = changes_df[
        changes_df[change_type_col].isin(
            ["INSERT", "UPDATE"]
        )
    ]

    latest_df = (
        changes_df[target_keys + [change_timestamp_col]]
        .groupby(
            target_keys,
            as_index=False,
            dropna=False,
        )
        .max()
        .merge(
          changes_df,
          how="inner",
          left_on=target_keys + [change_timestamp_col],
          right_on=target_keys + [change_timestamp_col],
        )
        .drop(columns=[change_type_col])
        .drop_duplicates(
          subset=target_keys, keep="first"
        )
        .rename(
            columns={
                change_timestamp_col: ingest_timestamp_col,
            }
        )
    )

    return latest_df
型変換
def cast_for_silver(
    df: bpd.DataFrame,
    target_schema: dict[str, str],
  ) -> bpd.DataFrame:
    """
    Silverテーブルのスキーマに合わせてDataFrameの列を型変換する。

    STRING列は変換せず、それ以外の列をSAFE_CASTする。
    戻り値の列順はtarget_schemaの定義順とする。

    Args:
        df: 型変換対象のBigQuery DataFramesのDataFrame。
        target_schema: 列名をキー、変換先のBigQueryデータ型を値とする辞書。

    Returns:
        Silverテーブルのスキーマに合わせて型変換したDataFrame。
    """

    result = df.copy()

    for column, target_type in target_schema.items():
        if column not in result.columns or target_type == "STRING":
            continue
        result[column] = bbq.sql_scalar(
            f"SAFE_CAST({{0}} AS {target_type})", columns=[result[column]]
        )

    output_cols = [
        column
        for column in target_schema
        if column in result.columns
    ]

    return result[output_cols]
MERGE
def merge_temp_into_target(
    client: bigquery.Client,
    target: str,
    source: str,
    temp_table: str,
    audit_cols: dict[str, str],
    end_ts: datetime | None = None,
    start_ts: datetime | None = None,
    tz: str | None = None,
    days_int: int | None = None,
  ) -> None:
    """
    ステージテーブルのデータをターゲットテーブルへMERGEする。

    主キーが一致し、ステージ側の取り込み日時が新しい場合は更新する。
    主キーが一致しない場合は新規挿入する。
    更新日時列にはMERGE実行時刻を設定する。

    Args:
        target: MERGE先の完全修飾テーブル名。
        source: MERGE元の完全修飾テーブル名。
        audit_cols: 監査列名の定義。

    Raises:
        ValueError: 主キーが不足している場合。
    """

    target_keys = get_primary_keys(
        client=client,
        table_full_name=target,
    )

    target_schema = get_table_schema(
        table_full_name=target,
    )
    target_columns = list(target_schema)

    ingest_timestamp_col = audit_cols.get(
        "ingest_timestamp_col"
    )
    update_timestamp_col = audit_cols.get(
        "update_timestamp_col"
    )

    temp_df = latest_by_keys(
        source=source,
        target_keys=target_keys,
        audit_cols=audit_cols,
        end_ts=end_ts,
        start_ts=start_ts,
        tz=tz,
        days_int=days_int,
    )

    temp_df = cast_for_silver(
        df=temp_df,
        target_schema=target_schema,
    )

    temp_columns = list(temp_df.columns)

    try:
      (
      temp_df.write
        .format("bigquery")
        .option("writeMethod", "direct")
        .option("createDisposition", "CREATE_IF_NEEDED")
        .mode("overwrite")
        .save(temp_table)
      )

      merge_condition = "\n      AND ".join(
          f"T.`{column}` = S.`{column}`"
          for column in target_keys
      )

      update_columns = [
          column 
          for column in target_columns
          if (
              column in temp_columns
              and column not in target_keys
              and column != update_timestamp_col
          )
      ]
      update_assignments = [
          f"T.`{column}` = S.`{column}`"
          for column in update_columns
      ]
      update_assignments.append(
          f"T.`{update_timestamp_col}` = CURRENT_TIMESTAMP()"
      )
      update_set = ",\n        ".join(
          update_assignments
      )

      insert_columns = [
          column
          for column in target_columns
          if (
              column in temp_columns
              or column == update_timestamp_col
          )
      ]
      insert_columns_sql = ", ".join(
          f"`{column}`"
          for column in insert_columns
      )
      insert_values_sql = ", ".join(
          (
              "CURRENT_TIMESTAMP()"
              if column == update_timestamp_col
              else f"S.`{column}`"
          )
          for column in insert_columns
      )

      merge_query = f"""
        MERGE `{target}` AS T
        USING `{temp_table}` AS S
          ON {merge_condition}
        WHEN MATCHED AND T.`{ingest_timestamp_col}` < S.`{ingest_timestamp_col}` THEN
          UPDATE SET
            {update_set}
        WHEN NOT MATCHED THEN
          INSERT ({insert_columns_sql})
          VALUES ({insert_values_sql})
      """

      client.query(merge_query).result()
      
    finally:
      client.delete_table(
          temp_table,
          not_found_ok=True,
      )

実行例

作成した関数に Bronze テーブル、Silver テーブル、一時テーブルの完全修飾名を渡して実行します。

dataset_id = "medallion_demo"

brz_table_name = 'product2__bronze'
slv_table_name = 'product2__silver'
temp_slv_table_name = '_temp_product2__silver'

audit_cols = {
    'ingest_timestamp_col': '_ingest_timestamp',
    'update_timestamp_col': '_update_timestamp'
}

brz_table_full_name = f"{PROJECT_ID}.{dataset_id}.{brz_table_name}"
slv_table_full_name = f"{PROJECT_ID}.{dataset_id}.{slv_table_name}"
temp_slv_table_full_name = f"{PROJECT_ID}.{dataset_id}.{temp_slv_table_name}"

merge_temp_into_target(
    client=bq_client,
    target=slv_table_full_name,
    source=brz_table_full_name,
    temp_table=temp_slv_table_full_name,
    audit_cols=audit_cols,
)

Spark in BigQuery

最後に、Spark in BigQuery を使って差分反映処理を実装します。

基本的な処理の流れは BigQuery DataFrames 版と同じです。CHANGES 関数で取得した変更履歴を Spark DataFrame として読み込み、主キーごとの最新レコードを抽出して型変換を行います。
変換結果を一時テーブルへ書き込んだ後、BigQuery の MERGE 文を実行して Silver テーブルへ反映します。

Spark ベースのデータ加工処理と BigQuery を組み合わせる場合や、既存のPySpark 処理へ Bronze から Silver への反映処理を組み込みたい場合に利用できる構成です。

処理の内容

処理は、主に次の流れで構成します。

  1. Spark セッションと BigQuery クライアントを作成する
  2. Silver テーブルの主キーとスキーマを取得する
  3. Bronze テーブルの変更履歴を Spark DataFrame として読み込む
  4. 主キーごとの最新レコードを取得する
  5. Silver テーブルのスキーマに合わせて型変換する
  6. 変換結果を一時テーブルへ書き込む
  7. BigQuery の MERGE 文で Silver テーブルへ反映する
  8. 処理後に一時テーブルを削除する
設定
from __future__ import annotations

from typing import Iterable
from datetime import datetime, timedelta, timezone
from zoneinfo import ZoneInfo
from uuid import uuid4

from google.cloud import bigquery
from google.cloud.dataproc_spark_connect import DataprocSparkSession
from google.cloud.dataproc_v1 import Session

from pyspark.sql import DataFrame
import pyspark.sql.connect.functions as F
from pyspark.sql.types import (
    BooleanType,
    DateType,
    DecimalType,
    DoubleType,
    IntegerType,
    LongType,
    StringType,
    StructField,
    StructType,
    TimestampType,
)

PROJECT_ID = "<PROJECT_ID>"
LOCATION = "asia-northeast1"

session = Session()
spark = (
    DataprocSparkSession.builder
    .appName("medallion-demo")
    .dataprocSessionConfig(session)
    .getOrCreate()
)

spark.conf.set("viewsEnabled", "true")
spark.conf.set("parentProject", PROJECT_ID)

bq_client = bigquery.Client(
    project=PROJECT_ID,
    location=LOCATION,
)
主キー取得
def get_primary_keys(
    client: bigquery.Client,
    table_full_name: str,
  ) -> list[str]:
    """
    BigQueryテーブルに定義された主キー列を取得する。

    Args:
        client: BigQueryクライアント
        table_full_name: project.dataset.table 形式のテーブル名。

    Returns:
        主キー列名のリスト。
    """

    table = client.get_table(table_full_name)
    constraints = table.table_constraints

    if constraints is None or constraints.primary_key is None:
        return []

    return list(constraints.primary_key.columns)
スキーマ取得
def get_table_schema(
    table_full_name: str
) -> dict[str, str]:
    """
    BigQueryコネクタが解釈したBigQueryテーブルの列名とSparkデータ型を取得する。

    Args:
        table_full_name: project.dataset.table 形式のテーブル名。

    Returns:
        列名をキー、BigQueryデータ型を値とする辞書。
    """

    schema = (
      spark.read
      .format("bigquery")
      .load(table_full_name)
      .schema
    )

    return {
        field.name: field.dataType.simpleString()
        for field in schema.fields
    }
主キーごとの最新レコードを取得
def latest_by_keys(
    source: str,
    target_keys: Iterable[str],
    audit_cols: dict[str, str],
    end_ts: datetime | None = None,
    start_ts: datetime | None = None,
    tz: str | None = None,
    days_int: int | None = None,
  ) -> DataFrame:
    """
    指定期間の変更履歴から、主キーごとの最新レコードを取得する。

    CHANGES関数が返すINSERTおよびUPDATEのレコードを対象とし、
    主キーごとに_CHANGE_TIMESTAMPが最も新しいレコードを残す。

    CHANGES関数の取得可能期間は1回につき最大1日であるため、
    指定期間を1日単位に分割して取得する。

    取得した_CHANGE_TIMESTAMPは、Silverテーブルで使用する
    取り込み日時列に設定する。

    Args:
        source:
            変更履歴を取得するBigQueryテーブル。
            project.dataset.table 形式で指定する。
        target_keys:
            最新レコードを判定する主キー列。
        audit_cols:
            監査列名の定義。
            ingest_timestamp_colキーを含む必要がある。
        end_ts:
            変更履歴の取得終了日時。指定時刻は含まれない。
            未指定の場合は現在日時。
        start_ts:
            変更履歴の取得開始日時。
            未指定の場合はend_tsからdays_int日前。
        tz:
            日時を解釈するタイムゾーン。
            未指定の場合はUTC。
        days_int:
            start_tsを省略した場合の取得日数。
            未指定の場合は7日。

    Returns:
        主キーごとの最新レコードを保持するBigQuery DataFrame。

    Raises:
        ValueError: 主キーが指定されていない場合。
    """

    target_keys = list(target_keys)

    if not target_keys:
        raise ValueError("keysには1つ以上の主キー列を指定してください。")

    ingest_timestamp_col = audit_cols.get("ingest_timestamp_col")

    days_int = 7 if days_int is None else days_int
    tz = "UTC" if tz is None else tz
    if end_ts is None:
        end_ts = datetime.now(ZoneInfo(tz))
    if start_ts is None:
        start_ts = end_ts - timedelta(days=days_int)

    source_columns = list(
        get_table_schema(
            table_full_name=source,)
    )

    # CHANGES関数の予約列には、通常の列名として使用できる別名を付ける。
    change_type_col = "change_type_internal"
    change_timestamp_col = "change_timestamp_internal"
    source_columns_sql = ",\n              ".join(
        f"`{column}`"
        for column in source_columns
    )

    change_queries: list[str] = []
    range_start = start_ts

    while range_start < end_ts:
        range_end = min(
            range_start + timedelta(days=1),
            end_ts,
        )
        range_start_literal = range_start.isoformat(
            timespec="microseconds"
        )
        range_end_literal = range_end.isoformat(
            timespec="microseconds"
        )
        change_queries.append(
            f"""
            SELECT
              {source_columns_sql},
              _CHANGE_TYPE AS `{change_type_col}`,
              _CHANGE_TIMESTAMP AS `{change_timestamp_col}`
            FROM CHANGES(
              TABLE `{source}`,
              TIMESTAMP('{range_start_literal}'),
              TIMESTAMP('{range_end_literal}')
            )
            """
        )
        range_start = range_end

    changes_query = "\nUNION ALL\n".join(change_queries)

    changes_df = (
      spark.read
      .format("bigquery")
      .option("query", changes_query)
      .load()
    )

    max_ts_df = (
        changes_df
        .filter(
          F.col(change_type_col).isin("INSERT", "UPDATE")
        )
        .select(
            *target_keys,
            F.col(change_timestamp_col),
        )
        .groupBy(*target_keys)
        .agg(
          F.max(F.col(change_timestamp_col)).alias(change_timestamp_col)
        )
    )

    join_columns = target_keys + [change_timestamp_col]

    latest_df = (
        max_ts_df
        .join(
          changes_df,
          on=join_columns,
          how="inner",
        )
        .drop(change_type_col)
        .dropDuplicates([*target_keys])
        .withColumnRenamed(
          change_timestamp_col,
          ingest_timestamp_col,
        )
    )

    return latest_df
型変換
def cast_for_silver(
    df: DataFrame,
    target_schema: dict[str | str],
) -> DataFrame:
    """
    入力DataFrameの列構成を維持し、ターゲットスキーマに
    定義された列のみ指定されたデータ型へ変換する。

    戻り値の列順は、入力DataFrameの列順とする。
    target_schemaに存在しない入力列は、型変換せずそのまま残す。

    Args:
        df:
            型変換対象のSpark DataFrame。
        target_schema:
            列名をキー、Sparkデータ型を値とする辞書。

    Returns:
        入力DataFrameの列構成に合わせて型変換したDataFrame。
    """

    source_columns = df.columns

    expressions = [
        F.col(f"`{column_name}`")
        .cast(target_schema[column_name])
        .alias(column_name)
        for column_name in source_columns
        if column_name in target_schema
    ]

    return df.select(*expressions)
MERGE
def merge_temp_into_target(
    client: bigquery.Client,
    target: str,
    source: str,
    temp_table: str,
    audit_cols: dict[str, str],
    end_ts: datetime | None = None,
    start_ts: datetime | None = None,
    tz: str | None = None,
    days_int: int | None = None,
  ) -> None:
    """
    Bronzeテーブルの変更履歴をSilverテーブルへ反映する。

    処理の流れは次のとおり。

    1. Silverテーブルに定義された主キーとスキーマを取得する。
    2. Bronzeテーブルの変更履歴から主キーごとの最新レコードを取得する。
    3. Silverテーブルのスキーマに合わせて型変換する。
    4. 変換結果を一時テーブルへ書き込む。
    5. 一時テーブルをSilverテーブルへMERGEする。
    6. 一時テーブルを削除する。

    主キーが一致し、ステージ側の取り込み日時が新しい場合は更新する。
    主キーが一致しない場合は新規挿入する。
    更新日時列にはMERGEの実行時刻を設定する。

    Args:
        client: 
            BigQueryクライアント
        target:
            MERGE先となるSilverテーブル。
        source:
            変更履歴を取得するBronzeテーブル。
        temp_table:
            MERGE元として使用する一時テーブル。
            実行のたびに置き換える。
        audit_cols:
            監査列名の定義。
            ingest_timestamp_colとupdate_timestamp_colを含む。
        end_ts:
            変更履歴の取得終了日時。
        start_ts:
            変更履歴の取得開始日時。
        tz:
            日時を解釈するタイムゾーン。
        days_int:
            start_tsを省略した場合の取得日数。
    """

    target_keys = get_primary_keys(
        client=client,
        table_full_name=target,
    )

    target_schema = get_table_schema(
        client=client,
        table_full_name=target
    )
    target_columns = list(target_schema)

    ingest_timestamp_col = audit_cols.get(
        "ingest_timestamp_col"
    )
    update_timestamp_col = audit_cols.get(
        "update_timestamp_col"
    )

    temp_df = latest_by_keys(
        client=client,
        source=source,
        target_keys=target_keys,
        audit_cols=audit_cols,
        end_ts=end_ts,
        start_ts=start_ts,
        tz=tz,
        days_int=days_int,
    )

    temp_df = cast_for_silver(
        df=temp_df,
        target_schema=target_schema,
    )

    temp_columns = list(temp_df.columns)

    try:
      # 戻り値には完全修飾テーブル名が返されるため、
      # MERGEではその値を使用する。
      temp_table_full_name = temp_df.to_gbq(
          temp_table,
          if_exists="replace",
          index=False,
      )

      merge_condition = "\n      AND ".join(
          f"T.`{key}` = S.`{key}`"
          for key in target_keys
      )

      update_columns = [
          column
          for column in target_columns
          if (
              column in temp_columns
              and column not in target_keys
              and column != update_timestamp_col
          )
      ]

      update_assignments = [
          f"T.`{column}` = S.`{column}`"
          for column in update_columns
      ]

      update_assignments.append(
          f"T.`{update_timestamp_col}` = CURRENT_TIMESTAMP()"
      )

      update_set = ",\n        ".join(
          update_assignments
      )

      insert_columns = [
          column
          for column in target_columns
          if (
              column in temp_columns
              or column == update_timestamp_col
          )
      ]

      insert_columns_sql = ", ".join(
          f"`{column}`"
          for column in insert_columns
      )

      insert_values_sql = ", ".join(
          (
              "CURRENT_TIMESTAMP()"
              if column == update_timestamp_col
              else f"S.`{column}`"
          )
          for column in insert_columns
      )

      merge_query = f"""
      MERGE `{target}` AS T
      USING `{temp_table_full_name}` AS S
        ON {merge_condition}
      WHEN MATCHED AND T.`{ingest_timestamp_col}` < S.`{ingest_timestamp_col}` THEN
        UPDATE SET
          {update_set}
      WHEN NOT MATCHED THEN
        INSERT ({insert_columns_sql})
        VALUES ({insert_values_sql})
      """

      client.query(merge_query).result()

    finally:
      client.delete_table(
        temp_table_full_name,
        not_found_ok=True,
      )

実行例

Bronze テーブル、Silver テーブル、一時テーブルの完全修飾名と監査列の設定を作成し、差分反映関数を実行します。

Spark から一時テーブルへ書き込む際の競合を避けるため、一時テーブル名には実行ごとに一意な値を使用します。

dataset_id = "medallion_demo"

brz_table_name = 'product2__bronze'
slv_table_name = 'product2__silver'
temp_slv_table_name = uuid4().hex

audit_cols = {
    'ingest_timestamp_col': '_ingest_timestamp',
    'update_timestamp_col': '_update_timestamp'
}

brz_table_full_name = f"{PROJECT_ID}.{dataset_id}.{brz_table_name}"
slv_table_full_name = f"{PROJECT_ID}.{dataset_id}.{slv_table_name}"
temp_slv_table_full_name = f"{PROJECT_ID}.{dataset_id}.{temp_slv_table_name}"

merge_temp_into_target(
    client=bq_client,
    target=slv_table_full_name,
    source=brz_table_full_name,
    temp_table=temp_slv_table_full_name,
    audit_cols=audit_cols,
)

一時テーブル名に関する注意点

固定の一時テーブル名を使用した場合、書き込みストリームの作成時にエラーが発生することがありました。
そのため、Spark in BigQuery 版では UUID を使って実行ごとに異なる一時テーブル名を生成しています。

Could not create write-stream after multiple retries
NOT_FOUND: Requested entity was not found
1
1
0

Register as a new user and use Qiita more conveniently

  1. You get articles that match your needs
  2. You can efficiently read back useful information
  3. You can use dark theme
What you can do with signing up
1
1

Delete article

Deleted articles cannot be recovered.

Draft of this article would be also deleted.

Are you sure you want to delete this article?