概要
incremental + insert_overwrite でコスト97%削減
materialized='table' は毎回フルリビルドで BQ 全スキャン。incremental + insert_overwrite に変えると直近パーティションのみ上書きされ、2TB → 60GB にスキャン量が激減する。
partition_by + cluster_by で剪定最大化
partition_by={'field':'order_date','data_type':'date','granularity':'day'} と cluster_by=['sku_id'] の組み合わせで、日付 + SKU の絞り込みクエリがほぼゼロコストになる。
COALESCE + SAFE_DIVIDE で NULL/ゼロ除算ガード
BQ では NULL * 数値 = NULL。COALESCE(quantity, 0) で集計狂いを防ぎ、SAFE_DIVIDE(sum, count) でゼロ除算を NULL に安全変換する。
dbt 1.8+ Unit Tests でリグレッション防止
unit_tests ブロックで arpu の NULL 保証・revenue 非負保証を CI に組み込む。--full-refresh 誤用による全リビルドをパイプラインレベルで防ぐ。
問題
ECサイトの MOps チームでは、Argo Workflows から BigQuery に毎時インクリメンタルロードするパイプラインを運用している。下記の dbt モデル群と BigQuery クエリには 6つの設計上の問題 が潜んでいる。問題点を全て洗い出し、dbt Core 1.8+・BigQuery のベストプラクティスに従って修正せよ。
シナリオ
毎時バッチで注文イベント(raw_orders)を受け取り、dbt で stg_orders → mart_daily_revenue へ変換し Looker Studio でダッシュボード表示する。現在のモデルでは1回の dbt run で BigQuery フルスキャン(約 2TB/run、$10/run)が発生しており、月間コスト $7,200 に達している。
制約・前提条件
- dbt Core 1.8+(
unique_idサロゲートキー、Unit Tests 機能使用可) - BigQuery パーティションは
order_date(DATE型)で DAY 分割 - クラスタリングキーは
sku_id - Incremental 戦略は
insert_overwrite(パーティション単位上書き) - 遅延データは最大 3 日なので
lookback_days = 3で処理 revenueはunit_price * quantityだが NULL の場合は 0 扱いarpuは購入者がいない場合は NULL(ゼロ除算を避ける)- dbt Unit Tests で
arpuの NULL 保証とrevenueの非負保証を追加すること
悪いコード (Before)
-- models/staging/stg_orders.sql
{{ config(
materialized='table' -- 問題①: 毎回フルリビルド
) }}
SELECT
order_id,
user_id,
sku_id,
quantity,
unit_price,
quantity * unit_price AS revenue, -- 問題②: NULLガードなし
status,
created_at, -- 問題③: パーティション列未定義
DATE(created_at) AS order_date
FROM {{ source('raw', 'orders') }}
WHERE status != 'cancelled'
---
-- models/mart/mart_daily_revenue.sql
{{ config(
materialized='table' -- 問題①: 毎回フルリビルド
) }}
SELECT
order_date,
sku_id,
SUM(revenue) AS total_revenue,
COUNT(DISTINCT user_id) AS unique_buyers,
SUM(quantity) AS total_qty,
SUM(revenue) / COUNT(DISTINCT user_id) AS arpu -- 問題④: ゼロ除算
FROM {{ ref('stg_orders') }}
GROUP BY 1, 2
from google.cloud import bigquery
def load_orders_to_bq(project_id: str, dataset: str):
client = bigquery.Client(project=project_id)
# 問題⑤: パーティション剪定なし(全件スキャン)
query = f"""
INSERT INTO `{project_id}.{dataset}.raw_orders`
SELECT *
FROM `{project_id}.{dataset}.orders_staging`
WHERE status IN ('purchased', 'shipped')
"""
client.query(query).result()
import subprocess
# 問題⑥: --full-refresh で毎回フルリビルド
subprocess.run([
"dbt", "run",
"--full-refresh",
"--select", "mart_daily_revenue"
])
materialized='table'(stg_orders・mart_daily_revenue 両方) — 毎回フルリビルドで BQ 全スキャン 2TB/run が発生。incremental + insert_overwrite で直近3日分のみ処理すれば 97% 削減quantity * unit_price の NULL ガードなし — どちらかが NULL だと revenue = NULL になり集計が狂う。COALESCE(quantity, 0) * COALESCE(unit_price, 0.0) で防ぐpartition_by 未設定(stg_orders) — パーティション定義がないと WHERE order_date >= ... でもフルスキャン。order_date DAY パーティション + sku_id クラスタリングを設定するarpu のゼロ除算ガードなし — unique_buyers = 0 になる場合に除算エラー。SAFE_DIVIDE で NULL を返すSELECT * で全件 INSERT。WHERE DATE(created_at) >= DATE_SUB(CURRENT_DATE(), INTERVAL 3 DAY) で差分のみ処理するdbt run --full-refresh を毎回実行 — incremental モデルに --full-refresh を渡すとフルリビルドになり問題①と同等。通常の dbt run + dbt test(Unit Tests)に変更するヒント(段階的開示)
ヒント1 — 方向性
materialized='table' は毎回全データを再ビルドするため BQ フルスキャンになる。materialized='incremental' に変えて is_incremental() マクロで差分のみ処理すること。パーティション列がないと WHERE order_date >= ... でも剪定が効かない。NULL の掛け算は NULL を返す(COALESCE が必要)。SUM / COUNT でゼロ除算が起きると BQ はエラーを返す。
ヒント2 — アプローチ
- 問題①:
materialized='incremental',incremental_strategy='insert_overwrite',partition_by={field:'order_date',data_type:'date',granularity:'day'}を設定 - 問題②:
COALESCE(quantity, 0) * COALESCE(unit_price, 0.0)で NULL を 0 に変換 - 問題③:
partition_by+cluster_by=['sku_id']を config に追加。is_incremental()ブロックでWHERE order_date >= DATE_SUB(CURRENT_DATE(), INTERVAL 3 DAY) - 問題④:
SAFE_DIVIDE(SUM(revenue), NULLIF(COUNT(DISTINCT user_id), 0))またはSAFE_DIVIDE(SUM(revenue), COUNT(DISTINCT user_id)) - 問題⑤: Loader クエリに
AND DATE(created_at) >= DATE_SUB(CURRENT_DATE(), INTERVAL 3 DAY)を追加 - 問題⑥:
--full-refreshを削除。dbt testを別途呼び出し Unit Tests を実行する
ヒント3 — コードの骨格
-- incremental model の骨格
{{ config(
materialized='incremental',
incremental_strategy='insert_overwrite',
partition_by={
'field': 'order_date',
'data_type': 'date',
'granularity': 'day'
},
cluster_by=['sku_id']
) }}
{% set LOOKBACK_DAYS = 3 %}
SELECT ...
{% if is_incremental() %}
WHERE order_date >= DATE_SUB(CURRENT_DATE(), INTERVAL {{ LOOKBACK_DAYS }} DAY)
{% endif %}
---
# dbt Unit Test 骨格(dbt 1.8+)
unit_tests:
- name: test_arpu_null_when_no_buyers
model: mart_daily_revenue
given:
- input: ref('stg_orders')
rows:
- {order_date: '2026-06-01', sku_id: 'SKU-001', revenue: 0.0, user_id: null, quantity: 0}
expect:
rows:
- {order_date: '2026-06-01', sku_id: 'SKU-001', arpu: null}
問題点分析(6点)
| # | 問題点 | 分類 | 改善方法 |
|---|---|---|---|
| 1 | materialized='table'(2モデル)でフルスキャン 2TB/run | コスト | incremental + insert_overwrite + is_incremental() で差分処理 |
| 2 | revenue = quantity * unit_price(NULL ガードなし) | データ品質 | COALESCE(quantity, 0) * COALESCE(unit_price, 0.0) |
| 3 | partition_by 未定義(stg_orders) | コスト/性能 | partition_by DAY + cluster_by=['sku_id'] |
| 4 | arpu = SUM / COUNT(ゼロ除算ガードなし) | データ品質 | SAFE_DIVIDE(SUM(revenue), COUNT(DISTINCT user_id)) |
| 5 | Loader クエリ全件 INSERT(パーティション剪定なし) | コスト | WHERE DATE(created_at) >= DATE_SUB(CURRENT_DATE(), INTERVAL 3 DAY) |
| 6 | dbt run --full-refresh を毎回実行(incremental の意味がない) | 信頼性 | 通常の dbt run + dbt test(Unit Tests) |
アーキテクチャ図 — Bad vs Good パイプライン比較(SVG)
模範解答
-- models/staging/stg_orders.sql
-- ECサイト注文ステージングモデル
-- 毎時インクリメンタルロード: order_date で DAY パーティション + sku_id クラスタリング
-- 遅延データ対応: lookback_days=3 で直近3日のパーティションを上書き
{{ config(
materialized='incremental', -- 修正①: フルリビルドを廃止
incremental_strategy='insert_overwrite', -- BQ ネイティブのパーティション上書き
partition_by={
'field': 'order_date', -- 修正③: DAY パーティション有効化
'data_type': 'date',
'granularity': 'day'
},
cluster_by=['sku_id'], -- 修正③: sku_id でクラスタリング
unique_key='order_id' -- 重複排除キー
) }}
{% set LOOKBACK_DAYS = 3 %} {# 遅延データの最大日数 #}
SELECT
order_id,
user_id,
sku_id,
-- 修正②: NULL ガード(quantity または unit_price が NULL の場合は 0.0 扱い)
COALESCE(quantity, 0) AS quantity,
COALESCE(unit_price, 0.0) AS unit_price,
COALESCE(quantity, 0) * COALESCE(unit_price, 0.0) AS revenue,
status,
created_at,
DATE(created_at) AS order_date -- パーティション列
FROM {{ source('raw', 'orders') }}
WHERE
status != 'cancelled'
-- 修正③: インクリメンタル実行時は直近 lookback_days のみスキャン
{% if is_incremental() %}
AND DATE(created_at) >= DATE_SUB(CURRENT_DATE(), INTERVAL {{ LOOKBACK_DAYS }} DAY)
{% endif %}
-- models/mart/mart_daily_revenue.sql
-- 日次売上マート: SKU×日付の集計
{{ config(
materialized='incremental', -- 修正①: テーブル全リビルドを廃止
incremental_strategy='insert_overwrite',
partition_by={
'field': 'order_date',
'data_type': 'date',
'granularity': 'day'
},
cluster_by=['sku_id']
) }}
{% set LOOKBACK_DAYS = 3 %}
SELECT
order_date,
sku_id,
SUM(revenue) AS total_revenue,
COUNT(DISTINCT user_id) AS unique_buyers,
SUM(quantity) AS total_qty,
-- 修正④: SAFE_DIVIDE でゼロ除算を NULL に変換
-- (購入者ゼロのセグメントでもクエリが落ちない)
SAFE_DIVIDE(
SUM(revenue),
NULLIF(COUNT(DISTINCT user_id), 0)
) AS arpu
FROM {{ ref('stg_orders') }}
{% if is_incremental() %}
WHERE order_date >= DATE_SUB(CURRENT_DATE(), INTERVAL {{ LOOKBACK_DAYS }} DAY)
{% endif %}
GROUP BY 1, 2
# models/mart/schema.yml
# dbt Core 1.8+ Unit Tests 定義
version: 2
models:
- name: mart_daily_revenue
description: "日次売上マート(SKU × 日付)"
columns:
- name: order_date
tests: [not_null]
- name: sku_id
tests: [not_null]
- name: total_revenue
tests: [not_null]
- name: arpu
description: "購入者がいない場合は NULL(SAFE_DIVIDE)"
unit_tests:
# テスト①: 購入者ゼロの場合に arpu が NULL になること(修正④保証)
- name: test_arpu_null_when_no_buyers
model: mart_daily_revenue
given:
- input: ref('stg_orders')
rows:
- {order_date: "2026-06-01", sku_id: "SKU-001",
revenue: 0.0, user_id: null, quantity: 0}
expect:
rows:
- {order_date: "2026-06-01", sku_id: "SKU-001", arpu: null}
# テスト②: revenue が正常に集計されること
- name: test_revenue_aggregation
model: mart_daily_revenue
given:
- input: ref('stg_orders')
rows:
- {order_date: "2026-06-01", sku_id: "SKU-002",
revenue: 1000.0, user_id: "U1", quantity: 2}
- {order_date: "2026-06-01", sku_id: "SKU-002",
revenue: 500.0, user_id: "U2", quantity: 1}
expect:
rows:
- {order_date: "2026-06-01", sku_id: "SKU-002",
total_revenue: 1500.0, unique_buyers: 2,
total_qty: 3, arpu: 750.0}
# テスト③: NULL quantity/unit_price が revenue=0 になること(修正②保証)
- name: test_revenue_non_negative_on_null_input
model: stg_orders
given:
- input: source('raw', 'orders')
rows:
- {order_id: "O1", user_id: "U1", sku_id: "SKU-003",
quantity: null, unit_price: 500.0,
status: "purchased", created_at: "2026-06-01 10:00:00 UTC"}
expect:
rows:
- {order_id: "O1", revenue: 0.0}
"""BigQuery インクリメンタルローダー。
Argo Workflows から毎時呼び出される。
直近 LOOKBACK_DAYS 日のデータのみを raw_orders に INSERT し、
dbt incremental run + unit tests を実行する。
"""
from __future__ import annotations
import argparse
import subprocess
import sys
from dataclasses import dataclass
from google.cloud import bigquery
# ---- 定数 ----------------------------------------------------------------
LOOKBACK_DAYS: int = 3 # 遅延データの最大日数(修正⑤⑥と合わせる)
DBT_SELECT_MODELS: str = "stg_orders mart_daily_revenue"
# -------------------------------------------------------------------------
@dataclass(frozen=True, slots=True)
class LoaderConfig:
"""ローダー設定値オブジェクト。
Attributes:
project_id: GCP プロジェクト ID
dataset: BigQuery データセット名
lookback_days: 処理対象の日数(遅延データ考慮)
"""
project_id: str
dataset: str
lookback_days: int = LOOKBACK_DAYS
def load_orders_incremental(config: LoaderConfig) -> None:
"""直近 lookback_days のオーダーデータを raw_orders にインクリメンタルロードする。
Args:
config: ローダー設定
Raises:
google.api_core.exceptions.GoogleAPIError: BQ クエリ失敗時
"""
client = bigquery.Client(project=config.project_id)
# 修正⑤: パーティション剪定を有効化(直近3日のみ INSERT)
# DATE(created_at) で order_date パーティションを剪定する
query = f"""
INSERT INTO `{config.project_id}.{config.dataset}.raw_orders`
SELECT *
FROM `{config.project_id}.{config.dataset}.orders_staging`
WHERE
status IN ('purchased', 'shipped')
AND DATE(created_at) >= DATE_SUB(
CURRENT_DATE(), INTERVAL {config.lookback_days} DAY
)
"""
job = client.query(query)
job.result() # 完了待機(例外は呼び出し元に伝播)
print(f"[INFO] BQ INSERT 完了: {job.num_dml_affected_rows} 行")
def run_dbt_incremental(select: str = DBT_SELECT_MODELS) -> None:
"""dbt incremental run + unit tests を実行する。
--full-refresh は渡さない(修正⑥)。
テスト失敗時はゼロ以外の終了コードで sys.exit する。
Args:
select: dbt --select に渡すモデル名(スペース区切り)
"""
# 修正⑥: --full-refresh を削除し通常の incremental run を実行
run_result = subprocess.run(
["dbt", "run", "--select", select],
capture_output=True, text=True,
)
if run_result.returncode != 0:
print(f"[ERROR] dbt run 失敗:\n{run_result.stderr}", file=sys.stderr)
sys.exit(run_result.returncode)
# Unit Tests は dbt test で別途実行(dbt 1.8+ unit_tests ブロック)
test_result = subprocess.run(
["dbt", "test", "--select", select],
capture_output=True, text=True,
)
if test_result.returncode != 0:
print(f"[ERROR] dbt test 失敗:\n{test_result.stderr}", file=sys.stderr)
sys.exit(test_result.returncode)
print("[INFO] dbt run + test 完了")
def main() -> None:
"""エントリポイント。"""
parser = argparse.ArgumentParser(description="BQ インクリメンタルローダー")
parser.add_argument("--project", required=True)
parser.add_argument("--dataset", required=True)
parser.add_argument("--lookback-days", type=int, default=LOOKBACK_DAYS)
args = parser.parse_args()
config = LoaderConfig(
project_id=args.project,
dataset=args.dataset,
lookback_days=args.lookback_days,
)
load_orders_incremental(config)
run_dbt_incremental()
if __name__ == "__main__":
main()
| 問題 | 修正内容 | 効果 |
|---|---|---|
| ① materialized='table'(2モデル) | incremental + insert_overwrite + is_incremental() | スキャン2TB→60GB(97%削減)$7,200→$216/月 |
| ② NULL ガードなし(revenue) | COALESCE(quantity, 0) * COALESCE(unit_price, 0.0) | NULL掛け算によるデータ欠損を防止 |
| ③ partition_by 未定義 | partition_by(order_date, DAY) + cluster_by(['sku_id']) | 日付+SKU絞り込みクエリがほぼゼロコスト |
| ④ ゼロ除算(arpu) | SAFE_DIVIDE(SUM(revenue), COUNT(DISTINCT user_id)) | 購入者ゼロセグメントでもクエリが安全に NULL を返す |
| ⑤ Loader 全件 INSERT | WHERE DATE(created_at) >= CURRENT_DATE - 3日 | Loader 側でもスキャン量を削減し二重コスト防止 |
| ⑥ --full-refresh 毎回実行 | dbt run(通常実行)+ dbt test(Unit Tests) | incremental を活かしつつリグレッションを CI で防止 |
コスト試算(修正前後)
| 項目 | 修正前 | 修正後 | 削減率 |
|---|---|---|---|
| スキャン量/run | 2 TB | 60 GB(3日分) | 97% |
| BQ コスト/run | $10.00 | $0.30 | 97% |
| 実行頻度 | 24回/日 | 24回/日 | — |
| 月間コスト | $7,200 | $216 | 97% |
| 月間削減額 | — | $6,984 | — |
partition_by + cluster_by の組み合わせによりダッシュボードクエリ(特定 SKU × 特定日付)のコストはさらに 90%以上削減される。ポイント解説
materialized='table' は毎回全テーブルを再ビルドするため BQ の全スキャン料金が発生する。insert_overwrite は is_incremental() が True のとき直近 3 日のパーティションのみを上書きするため、2TB → 約 60GB にスキャン量が激減する。初回・スキーマ変更時のみ dbt run --full-refresh を手動実行する運用が正しい。
BigQuery では
NULL * 数値 = NULL・NULL + 数値 = NULL となる。COALESCE(quantity, 0) * COALESCE(unit_price, 0.0) とすることで NULL を 0 として扱い、集計値が欠損する問題を防ぐ。NULL の意味が「未入力」ではなく「欠損」の場合は 0 への変換が適切であることを dbt の description に明記しておくと後任エンジニアが混乱しない。
パーティション定義がないと
WHERE order_date >= ... でもテーブル全スキャンになる。partition_by で DAY パーティションを設定し、cluster_by=['sku_id'] を加えることで WHERE order_date = '2026-06-03' AND sku_id = 'SKU-001' のクエリがほぼゼロコストになる。Looker Studio のダッシュボードが大量のリアルタイムクエリを発行する場合、この設定の有無でコストが 10〜100 倍変わる。
BigQuery の
SAFE_DIVIDE(x, y) は y=0 のとき NULL を返す。x / y は DivisionByZero エラー、x / NULLIF(y, 0) は NULL になるが、SAFE_DIVIDE が最も読みやすい dbt/BQ 標準パターン。Looker Studio で arpu = NULL のセルは空白表示になるため、ダッシュボード利用者への説明として column description に「購入者がいない場合は NULL」と記載する。
BQ Loader 側でも全件 INSERT していると、dbt の incremental 最適化が無意味になる。
WHERE DATE(created_at) >= DATE_SUB(CURRENT_DATE(), INTERVAL 3 DAY) で Loader とdbt の lookback_days を揃えることが重要。LOOKBACK_DAYS を定数化して両者で共有することでズレを防ぐ(dbt vars と Python 定数を同期させる CI チェックを追加するとより堅牢)。
dbt 1.8 の
unit_tests ブロックはモデルの SQL ロジックを入力データを固定して検証できる。--full-refresh を誤って本番に渡すと全リビルド料金が発生するため、CI でも dbt test(Unit Tests)で arpu の NULL 保証・revenue 非負保証を担保する。Argo Workflows のステップを dbt run → dbt test の2ステップに分け、dbt test 失敗時は Pub/Sub でアラートを飛ばす設計が MOps 標準パターン。
実務への応用
- Argo Workflows との連携:
bq_loader.pyを Argo WorkflowTemplate のscriptステップとして定義し、lookback_daysをワークフローパラメータ化する。失敗時はretryStrategyで最大 3 回リトライ、Dead Letter Queue(Pub/Sub)にエラー通知を送る - DataDog 可観測性:
job.num_dml_affected_rowsをカスタムメトリクス(custom_metrics.bq_loader.rows_inserted)として DataDog に送信し、異常値(極端に少ない挿入行数)をモニターで検知する。dbt のdbt_artifactsテーブルから実行時間・行数をメトリクスとして DataDog に送るdbt-datadogアダプターも活用できる - dbt Mesh(cross-project refs): MOps チームの
mart_daily_revenueを他チーム(マーケティング)がref('mops', 'mart_daily_revenue')で参照できるよう dbt Mesh + cross-project 参照を設定する。この際、ソースのpartition_by設定が参照先でも継承されることを確認する - BigQuery Materialized Views との使い分け:
mart_daily_revenueが Looker Studio から秒単位でクエリされる場合、dbt incremental の代わりに BigQuery Materialized Views(自動更新・クエリリライト)を検討する。ただし MV は JOINや HAVING に制約があるため、複雑な集計は dbt incremental の方が柔軟
今日のまとめ
incremental + insert_overwrite × BigQuery の partition_by + cluster_by は BQ コスト削減の最重要パターンであり(2TB→60GB, 97%削減)、COALESCE + SAFE_DIVIDE + dbt 1.8 Unit Tests の組み合わせでデータ品質とリグレッション防止を同時に担保する。特に
--full-refresh を誤って毎回実行すると incremental の恩恵がゼロになるため、Argo Workflows のパイプラインから --full-refresh フラグを除外し、初回・スキーマ変更時のみ手動実行する運用ルールを CI/CD ゲートで強制することが肝要。