はじめに
メダリオンアーキテクチャでは、ソースから取り込んだデータを 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 をすべて直接記述する方法と比べて、処理を関数単位に分割しやすく、複数のテーブルに適用できる共通処理として整理しやすい点が特徴です。
処理の内容
処理は、主に次の流れで構成します。
- Silver テーブルに定義された主キーを取得する
- Bronze テーブルと Silver テーブルのスキーマを取得する
- Bronze テーブルの変更履歴から主キーごとの最新レコードを取得する
- Silver テーブルのスキーマに合わせて型変換する
- 処理結果を一時テーブルへ書き込む
- 一時テーブルから Silver テーブルへ MERGE する
- 処理後に一時テーブルを削除する
以降では、それぞれの処理を関数に分けて実装しています。
設定
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 への反映処理を組み込みたい場合に利用できる構成です。
処理の内容
処理は、主に次の流れで構成します。
- Spark セッションと BigQuery クライアントを作成する
- Silver テーブルの主キーとスキーマを取得する
- Bronze テーブルの変更履歴を Spark DataFrame として読み込む
- 主キーごとの最新レコードを取得する
- Silver テーブルのスキーマに合わせて型変換する
- 変換結果を一時テーブルへ書き込む
- BigQuery の
MERGE文で Silver テーブルへ反映する - 処理後に一時テーブルを削除する
設定
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