C データエンジニアリング — dbt Unit Tests(1.8+)× BigQuery ARRAY_AGG × STRUCT × UNNEST × Materialized View × dbt Mesh cross-project ref × incremental insert_overwrite × var('lookback_days') × NULLIF ゼロ除算ガード(MOps 会員ランク×チャネル集計 Bad→Good)

2026-06-24 (Day 79) 水曜 C: データエンジニアリング ★★★★☆ BigQuery / dbt Core 1.8+ / Unit Tests / Mesh Materialized View / ARRAY_AGG / incremental insert_overwrite

概要

🧪

dbt Unit Tests(1.8+)でビジネスロジックを自動検証

dbt Core 1.8 から追加された Unit Tests 機能を使えば、given / expect 形式でインライン fixture を記述し、ランク変換ロジック(bronze/silver/gold/platinum → rank_weight)や NULLIF ゼロ除算ガードを SQL だけでテストできる。dbt test --select mart_member_channel_monthly を CI に組み込むことでリグレッションを防ぐ。

🔢

ARRAY_AGG × STRUCT × UNNEST で Python ループを廃絶

会員ランク × チャネル = 12 組み合わせを Python for ループで BigQuery API に 12 回問い合わせるアンチパターンは、ARRAY_AGG(STRUCT(...)) でネスト集約し UNNEST で平坦化することで 1 クエリに統合できる。BigQuery の分散処理エンジンに委譲することで処理時間が数分から数十秒に短縮される。

💎

BigQuery Materialized View で下流スキャンコストをゼロに

下流 analytics プロジェクトが同じ集計クエリを毎回フルスキャンする構造は、BigQuery Materialized View(enable_refresh: true, refresh_interval_minutes: 60)で解消できる。差分更新のみ課金されるため月額 $200+ → ~$10 に削減。dbt Mesh cross-project ref 経由で安定した公開 API として提供できる。

🔄

incremental insert_overwrite で遅延データを自動補正

incremental_strategy = 'insert_overwrite' + partition_by = {field: 'report_month'} を組み合わせると、月次パーティション単位で上書きが走る。var('lookback_days', 90) のルックバック窓を設けることで Pub/Sub 遅延によって翌日到着する注文データを自動的に取り込み、集計の整合性を保つ。

問題

ECサイトの MOps チームでは、BigQuery × dbt Core 1.8+ でキャンペーン別 会員ランク × 購買チャネル クロス集計レポート を生成している。以下の Python スクリプト・dbt モデルには 7つの設計上の問題 が潜んでいる。問題点を全て洗い出し、BigQuery ARRAY_AGG × STRUCT × UNNEST・dbt Unit Tests・dbt Mesh cross-project ref・dbt incremental insert_overwrite・BigQuery Materialized View を活用した Bad→Good リファクタリングを行え。

制約・前提条件

  • dbt Core 1.8+(Unit Tests 機能・Mesh cross-project ref 対応)
  • BigQuery Standard SQL(GA 版機能のみ)
  • テーブル: events.member_actions(partition: action_date DATE, cluster: member_rank), orders.order_items(partition: order_date DATE, cluster: channel
  • 集計軸: 会員ランク(bronze/silver/gold/platinum)× 購買チャネル(web/app/store)× 月
  • Argo Workflows CronWorkflow で毎朝 JST 06:00 実行
  • 下流モデル(analytics プロジェクト)が dbt Mesh cross-project ref で参照
期待する回答形式: 問題点の列挙(番号付き)+ 改善後 dbt モデル(SQL)+ dbt Unit Test YAML + sources.yml + dbt_project.yml(vars)+ 設計意図の説明

悪いコード (Before)

このコード・モデルには 7つのデータエンジニアリング設計の問題 が隠れています。
bad_member_channel.py — SELECT * / Python ループ / マジックナンバー / dbt ハードコード
from google.cloud import bigquery
import pandas as pd

client = bigquery.Client()

RANKS   = ["bronze", "silver", "gold", "platinum"]
CHANNELS = ["web", "app", "store"]

def run_cross_agg():
    results = []
    for rank in RANKS:
        for channel in CHANNELS:
            # 問題①: SELECT * → 40カラム全スキャン(月額 $180 無駄)
            # 問題②: ループ × BigQuery API → 12回クエリ(N+1問題)
            # 問題③: マジックナンバー 90 がハードコード
            q = f"""
                SELECT *
                FROM `my-project.events.member_actions` AS m
                JOIN `my-project.orders.order_items` AS o
                  ON m.user_id = o.user_id
                WHERE m.member_rank = '{rank}'
                  AND o.channel = '{channel}'
                  AND o.order_date >= DATE_SUB(CURRENT_DATE(), INTERVAL 90 DAY)
                GROUP BY 1, 2, 3
            """
            # 問題④: GROUP BY 1,2,3 数字インデックス
            df = client.query(q).to_dataframe()
            results.append(df)

    final = pd.concat(results)
    # 問題⑦: 下流が毎回この結果をフルスキャン(Materialized View なし)
    final.to_gbq("mart.member_channel_monthly", if_exists="replace")

# ── dbt モデル ──
# models/mart/mart_member_channel_monthly.sql
# 問題⑤: テーブル名ハードコード(ref()/source()を使わない)
# SELECT member_rank, channel, SUM(revenue)
# FROM `my-project.events.member_actions`   ← ハードコード
# GROUP BY 1, 2
# 問題⑥: dbt Unit Tests 未設定(ランク変換ロジックの回帰テストなし)
問題点サマリー(7点)
1SELECT * フルカラムスキャン — 40カラム全取得で月額 $180。必要5カラムのみで $22.5(87.5%削減)
2Python ループ N+1 BigQuery API — 12組み合わせ × BigQuery クエリ 12回。ARRAY_AGG × UNNEST で 1クエリ化
3マジックナンバー 90 の散在 — SQL / Python に 3 箇所ハードコード。{{ var('lookback_days', 90) }} で一元管理
4GROUP BY 1,2,3 数字インデックス — カラム順変更でサイレントに集計軸が変わる。カラム名明示で保守性向上
5dbt ref() / source() 未使用 — 依存グラフが機能せず dbt Mesh cross-project ref が断裂する
6dbt Unit Tests 未設定 — rank_weight CASE 式の回帰テストなし。dbt 1.8+ Unit Tests で自動保護
7Materialized View 未使用 — 下流 analytics が毎回フルスキャン(月額 $200+)→ Materialized View で ~$10 に削減

ヒント(段階的開示)

ヒント1 — 方向性
現状コードの最大問題は「Python で会員ランク × チャネルの組み合わせをループ処理」している点。BigQuery は ARRAY_AGG × STRUCT で行方向にネストしたデータを 1 クエリで生成でき、UNNEST で平坦化できる。また dbt Unit Tests(1.8+)を使えばビジネスロジックの回帰テストを SQL fixture だけで記述できる。
ヒント2 — アプローチ
  • 問題①: SELECT * → 必要カラム 5 つのみ明示(スキャン 87.5% 削減)
  • 問題②: Python for ループ × BigQuery API → BigQuery JOIN で一括処理(1 クエリ化)
  • 問題③: マジックナンバー 90{{ var('lookback_days', 90) }}(dbt_project.yml で管理)
  • 問題④: GROUP BY 1,2,3GROUP BY member_rank, channel, report_month(カラム名明示)
  • 問題⑤: ハードコード → {{ source('events', 'member_actions') }} / {{ source('orders', 'order_items') }}
  • 問題⑥: Unit Tests なし → tests/unit/test_mart_member_channel_monthly.ymlgiven / expect を記述
  • 問題⑦: 下流フルスキャン → BigQuery Materialized View(enable_refresh: true, refresh_interval_minutes: 60
ヒント3 — コードの骨格
-- BigQuery JOIN で一括処理 + var() で lookback 管理(骨格)
WITH
member_actions AS (
  SELECT user_id, member_rank, action_date
  FROM {{ source('events', 'member_actions') }}   -- 修正⑤
  {% if is_incremental() %}
    WHERE action_date >= DATE_SUB(
      CURRENT_DATE('Asia/Tokyo'),
      INTERVAL {{ var('lookback_days', 90) }} + 1 DAY  -- 修正③
    )
  {% endif %}
),
order_items AS (
  SELECT user_id, order_id, channel,
    DATE_TRUNC(DATE(order_timestamp, 'Asia/Tokyo'), MONTH) AS report_month,
    revenue
  FROM {{ source('orders', 'order_items') }}
  {% if is_incremental() %}
    WHERE order_date >= DATE_SUB(CURRENT_DATE('Asia/Tokyo'),
      INTERVAL {{ var('lookback_days', 90) }} + 1 DAY)
  {% endif %}
),
aggregated AS (
  SELECT
    m.member_rank,       -- 修正④: カラム名明示
    o.channel,
    o.report_month,
    COUNT(DISTINCT o.order_id) AS order_count,
    ROUND(SUM(o.revenue), 2)   AS total_revenue,
    CASE m.member_rank
      WHEN 'platinum' THEN 4 WHEN 'gold' THEN 3
      WHEN 'silver'   THEN 2 WHEN 'bronze' THEN 1
      ELSE 0 END         AS rank_weight
  FROM member_actions AS m
  INNER JOIN order_items AS o ON m.user_id = o.user_id
  GROUP BY member_rank, channel, report_month  -- 修正④
)
SELECT * FROM aggregated

dbt Unit Test YAML 骨格:

unit_tests:
  - name: test_rank_weight_platinum
    model: mart_member_channel_monthly
    given:
      - input: source('events', 'member_actions')
        rows:
          - {user_id: "u1", member_rank: "platinum", action_date: "2026-01-15"}
      - input: source('orders', 'order_items')
        rows:
          - {user_id: "u1", order_id: "o1", channel: "web",
             order_timestamp: "2026-01-20T10:00:00Z", revenue: 10000.0}
    expect:
      rows:
        - {member_rank: "platinum", rank_weight: 4}

問題点分析(7点)

#問題点分類改善方法
1SELECT * フルカラムスキャン(40カラム)コスト必要5カラムのみ明示(87.5%削減)
2Python ループ N+1 BigQuery API(12回)パフォーマンスBigQuery JOIN 1本に集約
3マジックナンバー 90 の散在(3箇所)保守性var('lookback_days', 90) で一元管理
4GROUP BY 1,2,3 数字インデックス保守性カラム名明示(member_rank, channel, report_month
5テーブル名ハードコード(ref()/source() 未使用)依存管理{{ source('events', 'member_actions') }}
6dbt Unit Tests 未設定品質1.8+ Unit Tests YAML で rank_weight / ゼロ除算を自動検証
7BigQuery Materialized View 未使用コストenable_refresh: true, refresh_interval_minutes: 60

設計図 — Bad vs Good データフロー

Bad(変更前)— Python ループ + ハードコード BigQuery: member_actions / order_items ① SELECT * → 40カラム全スキャン(月額 $180) ⚠️ テーブル名ハードコード(source() なし)→ dbt 依存グラフ機能しない ⑤ dbt Mesh cross-project ref が断裂する Python ループ(Argo Workflows Pod) ② for rank in RANKS: for channel in CHANNELS: → 12回クエリ(N+1) ⚠️ 12組み合わせ × BigQuery API 呼び出し → 処理時間: 数分 ③ INTERVAL 90 DAY — マジックナンバーが SQL と Python に散在(3箇所) ④ GROUP BY 1, 2, 3 — 数字インデックス(カラム追加でサイレント集計軸変化) ⑥ dbt Unit Tests なし → rank_weight CASE 式の回帰リスクあり ⚠️ 不明ランクが来ても検知できず(ELSE 句なし) mart.member_channel_monthly(to_gbq) Pandas concat → to_gbq → replace 全件洗替(増分処理なし) ⚠️ 遅延到着データが翌日に補正されない 下流 analytics プロジェクト(毎回フルスキャン) ⑦ Materialized View なし → SELECT SUM(revenue)... で毎回全件スキャン ⚠️ 月額 $200+ / Looker Studio ダッシュボード表示 10秒以上 dbt Mesh cross-project ref が断裂 → テーブル名ハードコードで参照 Good(変更後)— BigQuery SQL + dbt Unit Tests BigQuery: {{ source('events/orders', ...) }} 修正①: user_id, member_rank, action_date, channel, revenue の5カラムのみ(月額 $22.5) 修正⑤: dbt source() で依存グラフ管理 / freshness warn 2h / error 6h dbt Mesh public model として access: public + latest_version: 1 で公開 パーティション剪定: action_date / order_date フィルタで incremental スキャン削減 BigQuery JOIN(1クエリ化) + incremental insert_overwrite 修正②: Python ループ廃止 → INNER JOIN で BigQuery 分散処理に委譲(処理: 数十秒) 修正③: INTERVAL {{ var('lookback_days', 90) }} DAY — dbt_project.yml で一元管理(本番:90, CI:7) 修正④: GROUP BY member_rank, channel, report_month — カラム名明示 rank_weight: CASE member_rank WHEN 'platinum' THEN 4 ... ELSE 0 END(未知ランクガード) avg_order_value: SUM(revenue) / NULLIF(COUNT(DISTINCT order_id), 0) — ゼロ除算ガード incremental_strategy: insert_overwrite + partition_by: report_month で月次単位上書き lookback_days ルックバックで Pub/Sub 遅延による翌日到着データを自動補正 dbt Unit Tests(1.8+)— 自動回帰テスト 修正⑥: test_rank_weight_platinum / test_rank_weight_unknown / test_zero_revenue_excluded given: インライン fixture → expect: 期待値 → dbt test --select mart_member_channel_monthly CI/CD(Cloud Build / GitHub Actions)に組み込み → マージ前に自動検証 accepted_values テスト: member_rank ∈ {bronze,silver,gold,platinum}, channel ∈ {web,app,store} mart.member_channel_monthly(正確・低コスト・テスト保護) ランク × チャネル × 月次集計 / スキャン87.5%削減 / 遅延データ自動補正済み access: public + dbt Mesh cross-project ref で analytics プロジェクトへ安定公開 BigQuery Materialized View(修正⑦) 修正⑦: enable_refresh: true, refresh_interval_minutes: 60 — 差分更新のみ課金 月額 $200+ → ~$10 / Looker Studio 表示: 10秒 → <1秒(キャッシュヒット) dbt Mesh cross-project ref: {{ ref('analytics', 'mv_member_channel_monthly', v=1) }} DataDog: dbt source freshness metric → Slack アラート(2h warn / 6h error) 修正

模範解答

mart_member_channel_monthly.sql — 改善後 dbt モデル

-- models/mart/mart_member_channel_monthly.sql
-- dbt Core 1.8+ | BigQuery Standard SQL
-- 集計軸: 会員ランク × 購買チャネル × 月次
{{
  config(
    materialized         = 'incremental',
    partition_by         = {'field': 'report_month', 'data_type': 'date'},
    cluster_by           = ['member_rank', 'channel'],
    on_schema_change     = 'fail',
    incremental_strategy = 'insert_overwrite',   -- 月次パーティション単位で上書き
    -- dbt Mesh: 下流 analytics プロジェクトへ公開
    meta                 = {'access': 'public', 'version': 1},
  )
}}

WITH
-- 修正①: 必要カラム 5 つのみ(スキャン 87.5% 削減)
-- 修正⑤: {{ source() }} で依存グラフを明示
member_actions AS (
  SELECT
    user_id,
    member_rank    -- bronze / silver / gold / platinum
  FROM {{ source('events', 'member_actions') }}     -- 修正⑤
  {% if is_incremental() %}
    -- インクリメンタル: lookback_days 日分のルックバックで遅延データも補正
    WHERE action_date >= DATE_SUB(
      CURRENT_DATE('Asia/Tokyo'),
      INTERVAL {{ var('lookback_days', 90) }} + 1 DAY  -- 修正③
    )
  {% endif %}
),

order_items AS (
  SELECT
    user_id,
    order_id,
    channel,                                      -- web / app / store
    DATE_TRUNC(
      DATE(order_timestamp, 'Asia/Tokyo'), MONTH
    ) AS report_month,                            -- JST 月次集計
    revenue
  FROM {{ source('orders', 'order_items') }}       -- 修正⑤
  {% if is_incremental() %}
    WHERE order_date >= DATE_SUB(
      CURRENT_DATE('Asia/Tokyo'),
      INTERVAL {{ var('lookback_days', 90) }} + 1 DAY  -- 修正③
    )
  {% endif %}
),

-- 修正②: Python ループ廃止 → BigQuery JOIN で 1 クエリ一括処理
-- 修正③: var('lookback_days', 90) で管理(dbt_project.yml: 本番=90, CI=7)
joined AS (
  SELECT
    m.member_rank,
    o.channel,
    o.report_month,
    o.order_id,
    o.revenue
  FROM member_actions AS m
  INNER JOIN order_items AS o
    ON m.user_id = o.user_id
    AND o.report_month >= DATE_TRUNC(
          DATE_SUB(
            CURRENT_DATE('Asia/Tokyo'),
            INTERVAL {{ var('lookback_days', 90) }} DAY  -- 修正③
          ),
          MONTH
        )
),

aggregated AS (
  SELECT
    -- 修正④: GROUP BY 1,2,3 廃止 → カラム名明示
    member_rank,
    channel,
    report_month,
    COUNT(DISTINCT order_id) AS order_count,
    ROUND(SUM(revenue), 2)   AS total_revenue,
    -- ゼロ除算ガード: order_count=0 のセルは NULL(COALESCE で 0 に変換可)
    ROUND(
      SUM(revenue) / NULLIF(COUNT(DISTINCT order_id), 0), 2
    )                        AS avg_order_value,
    -- 会員ランク重みスコア(下流 BI 用)
    -- 修正⑥: dbt Unit Tests でこの CASE 式を自動検証
    CASE member_rank
      WHEN 'platinum' THEN 4  -- 最上位: ポイント還元率 3%
      WHEN 'gold'     THEN 3  -- 上位: ポイント還元率 2%
      WHEN 'silver'   THEN 2  -- 中位: ポイント還元率 1.5%
      WHEN 'bronze'   THEN 1  -- 通常: ポイント還元率 1%
      ELSE 0                  -- 未知ランク: ゼロで安全処理(guard)
    END                      AS rank_weight
  FROM joined
  -- 修正④: カラム名で GROUP BY(数字インデックス廃止)
  GROUP BY member_rank, channel, report_month
)

SELECT * FROM aggregated

tests/unit/test_mart_member_channel_monthly.yml — dbt Unit Tests(修正⑥)

# tests/unit/test_mart_member_channel_monthly.yml
# dbt Core 1.8+ Unit Tests
# 実行: dbt test --select mart_member_channel_monthly

unit_tests:
  # ── ランク変換ロジックの検証 ──────────────────────────────
  - name: test_rank_weight_platinum
    description: "platinum 会員は rank_weight=4 になること(最上位ランク)"
    model: mart_member_channel_monthly
    given:
      - input: source('events', 'member_actions')
        rows:
          - {user_id: "u1", member_rank: "platinum", action_date: "2026-01-15"}
      - input: source('orders', 'order_items')
        rows:
          - {user_id: "u1", order_id: "o1", channel: "web",
             order_timestamp: "2026-01-20T10:00:00Z", revenue: 10000.0}
    expect:
      rows:
        - {member_rank: "platinum", channel: "web", order_count: 1,
           total_revenue: 10000.0, rank_weight: 4}

  - name: test_rank_weight_all_tiers
    description: "全ランクの rank_weight が正しくマッピングされること"
    model: mart_member_channel_monthly
    given:
      - input: source('events', 'member_actions')
        rows:
          - {user_id: "u1", member_rank: "gold",     action_date: "2026-01-15"}
          - {user_id: "u2", member_rank: "silver",   action_date: "2026-01-15"}
          - {user_id: "u3", member_rank: "bronze",   action_date: "2026-01-15"}
      - input: source('orders', 'order_items')
        rows:
          - {user_id: "u1", order_id: "o1", channel: "app",
             order_timestamp: "2026-01-20T10:00:00Z", revenue: 8000.0}
          - {user_id: "u2", order_id: "o2", channel: "app",
             order_timestamp: "2026-01-20T11:00:00Z", revenue: 5000.0}
          - {user_id: "u3", order_id: "o3", channel: "app",
             order_timestamp: "2026-01-20T12:00:00Z", revenue: 3000.0}
    expect:
      rows:
        - {member_rank: "gold",   rank_weight: 3}
        - {member_rank: "silver", rank_weight: 2}
        - {member_rank: "bronze", rank_weight: 1}

  - name: test_rank_weight_unknown_guard
    description: "未知のランクは rank_weight=0 になること(ELSE 0 ガード)"
    model: mart_member_channel_monthly
    given:
      - input: source('events', 'member_actions')
        rows:
          - {user_id: "u9", member_rank: "vip_unknown", action_date: "2026-01-15"}
      - input: source('orders', 'order_items')
        rows:
          - {user_id: "u9", order_id: "o9", channel: "store",
             order_timestamp: "2026-01-20T10:00:00Z", revenue: 50000.0}
    expect:
      rows:
        - {member_rank: "vip_unknown", rank_weight: 0}

  # ── ゼロ除算ガードの検証 ──────────────────────────────────
  - name: test_zero_revenue_avg_order_value
    description: "revenue=0 の注文は avg_order_value=0.0 になること(NULLIF ゼロ除算ガード)"
    model: mart_member_channel_monthly
    given:
      - input: source('events', 'member_actions')
        rows:
          - {user_id: "u4", member_rank: "bronze", action_date: "2026-01-15"}
      - input: source('orders', 'order_items')
        rows:
          - {user_id: "u4", order_id: "o4", channel: "store",
             order_timestamp: "2026-01-20T10:00:00Z", revenue: 0.0}
    expect:
      rows:
        - {member_rank: "bronze", avg_order_value: 0.0}

  # ── マルチチャネル集計の検証 ──────────────────────────────
  - name: test_multichannel_aggregation
    description: "同一ユーザーが複数チャネルで購買した場合、チャネル別に集計されること"
    model: mart_member_channel_monthly
    given:
      - input: source('events', 'member_actions')
        rows:
          - {user_id: "u5", member_rank: "gold", action_date: "2026-01-15"}
      - input: source('orders', 'order_items')
        rows:
          - {user_id: "u5", order_id: "w1", channel: "web",
             order_timestamp: "2026-01-20T10:00:00Z", revenue: 12000.0}
          - {user_id: "u5", order_id: "a1", channel: "app",
             order_timestamp: "2026-01-21T10:00:00Z", revenue: 8000.0}
    expect:
      rows:
        - {member_rank: "gold", channel: "web", order_count: 1, total_revenue: 12000.0}
        - {member_rank: "gold", channel: "app", order_count: 1, total_revenue: 8000.0}

models/sources.yml — freshness 設定(修正⑤ + 可観測性)

# models/sources.yml
version: 2

sources:
  - name: events
    database: my-gcp-project
    schema: events
    tables:
      - name: member_actions
        description: "会員アクションイベント(Pub/Sub → BQ streaming insert)"
        # freshness 監視: Pub/Sub → BQ パイプライン遅延を自動検知
        loaded_at_field: _ingested_at   # streaming insert のシステム列
        freshness:
          warn_after:  {count: 2, period: hour}   # 2時間以上新規データなし → warn
          error_after: {count: 6, period: hour}   # 6時間以上 → error (Slack アラート)
        columns:
          - name: user_id
            tests: [not_null]
          - name: member_rank
            tests:
              - not_null
              - accepted_values:
                  values: ["bronze", "silver", "gold", "platinum"]
          - name: action_date
            tests: [not_null]

  - name: orders
    database: my-gcp-project
    schema: orders
    tables:
      - name: order_items
        description: "注文明細(チャネル: web/app/store)"
        loaded_at_field: _ingested_at
        freshness:
          warn_after:  {count: 2, period: hour}
          error_after: {count: 6, period: hour}
        columns:
          - name: order_id
            tests: [not_null, unique]
          - name: channel
            tests:
              - not_null
              - accepted_values:
                  values: ["web", "app", "store"]
          - name: revenue
            tests: [not_null]

BigQuery Materialized View — 下流スキャンコスト排除(修正⑦)

-- analytics-project: 下流から参照される Materialized View
-- Terraform で管理する場合は google_bigquery_table リソース
-- dbt Mesh cross-project ref: {{ ref('my_project', 'mart_member_channel_monthly', v=1) }}

CREATE MATERIALIZED VIEW IF NOT EXISTS
  `analytics-project.mart.mv_member_channel_monthly`
OPTIONS (
  -- 1時間ごとに差分更新(更新分のみ課金)
  enable_refresh           = true,
  refresh_interval_minutes = 60,
  -- 最後の更新時刻を DataDog で監視可能
  description              = 'member_rank x channel 月次集計 | dbt mart参照 | updated: 60min interval'
) AS
SELECT
  member_rank,
  channel,
  report_month,
  SUM(order_count)                           AS order_count,
  ROUND(SUM(total_revenue), 2)               AS total_revenue,
  ROUND(SUM(total_revenue)
    / NULLIF(SUM(order_count), 0), 2)        AS avg_order_value,
  MAX(rank_weight)                           AS rank_weight
FROM `my-gcp-project.mart.mart_member_channel_monthly`
GROUP BY member_rank, channel, report_month;

-- ── Terraform での管理 ───────────────────────────────
-- resource "google_bigquery_table" "mv_member_channel_monthly" {
--   dataset_id = "mart"
--   table_id   = "mv_member_channel_monthly"
--   project    = var.analytics_project_id
--
--   materialized_view {
--     query                            = <<-SQL
--       SELECT member_rank, channel, report_month,
--              SUM(order_count) AS order_count, ...
--       FROM `${var.project_id}.mart.mart_member_channel_monthly`
--       GROUP BY member_rank, channel, report_month
--     SQL
--     enable_refresh           = true
--     refresh_interval_ms      = 3600000   # 60 分
--   }
-- }

dbt_project.yml — vars + incremental 設定

# dbt_project.yml(抜粋)
name: my_project
version: "1.0.0"
config-version: 2

# 修正③: マジックナンバー一元管理
# 本番: lookback_days=90 / CI: dbt run --vars '{"lookback_days": 7}'
vars:
  lookback_days: 90            # attribution ルックバック日数(Pub/Sub 遅延対応)

models:
  my_project:
    mart:
      +materialized: incremental
      +on_schema_change: fail
      +partition_by:
        field: report_month
        data_type: date
      +cluster_by: ["member_rank", "channel"]
      +incremental_strategy: insert_overwrite

# dbt Mesh: public model として公開
# models/mart/mart_member_channel_monthly.yml で access: public を設定
# 下流から: {{ ref('my_project', 'mart_member_channel_monthly', v=1) }}

Argo Workflows CronWorkflow(JST 06:00 実行)

# dbt-member-channel-cronworkflow.yaml
apiVersion: argoproj.io/v1alpha1
kind: CronWorkflow
metadata:
  name: dbt-member-channel-monthly
  namespace: mops
spec:
  # JST 06:00 = UTC 21:00(前日)
  schedule: "0 21 * * *"
  timezone: "Asia/Tokyo"
  workflowSpec:
    entrypoint: dbt-pipeline
    templates:
      - name: dbt-pipeline
        steps:
          - - name: source-freshness
              template: dbt-cmd
              arguments:
                parameters:
                  - {name: cmd, value: "dbt source freshness --output json"}
          - - name: run-mart
              template: dbt-cmd
              arguments:
                parameters:
                  - {name: cmd, value: "dbt run --select mart_member_channel_monthly"}
          - - name: test-mart
              template: dbt-cmd
              arguments:
                parameters:
                  - {name: cmd, value: "dbt test --select mart_member_channel_monthly"}
      - name: dbt-cmd
        inputs:
          parameters:
            - name: cmd
        container:
          image: ghcr.io/dbt-labs/dbt-bigquery:1.8
          command: [sh, -c]
          args: ["{{inputs.parameters.cmd}}"]
          env:
            - name: DBT_PROFILES_DIR
              value: /workspace
          resources:
            requests:
              memory: "512Mi"
              cpu: "500m"

ポイント解説

  1. dbt Unit Tests(1.8+)の given/expect 構造: given セクションでテスト用インライン fixture を定義し、expect で期待値を宣言する。従来の schema.yml テスト(not_null / accepted_values)では検証できなかった「CASE 式のビジネスロジック」や「ゼロ除算ガード」を SQL fixture だけで保護できる。dbt test --select mart_member_channel_monthly を CI に組み込むことでマージ前に自動検証が走る。
  2. NULLIF ゼロ除算ガード: SUM(revenue) / NULLIF(COUNT(DISTINCT order_id), 0) は分母が 0 の場合に NULL を返す。BigQuery はゼロ除算で NULL(エラーなし)を返すが、明示的な NULLIF は意図を示すドキュメントとしての役割も果たす。ダッシュボード側では COALESCE(avg_order_value, 0) で 0 表示に変換する。
  3. incremental insert_overwrite + lookback: partition_by = report_month + insert_overwrite の組み合わせは「当月パーティションを毎実行で上書き」する。lookback_days ルックバック窓(90日)を設けることで Pub/Sub → BQ 間の遅延(最大数時間)による翌日到着データを自動的に取り込み、月次集計の整合性を保つ。本番: 90日、CI: 7日(dbt run --vars '{"lookback_days": 7}')で使い分け。
  4. BigQuery Materialized View のコスト効果: Materialized View は定義クエリの結果をキャッシュし、refresh_interval_minutes: 60 で差分のみ更新する。下流 analytics プロジェクトからの参照はキャッシュヒット時にスキャンゼロとなるため月額 $200+ → ~$10 に削減。Looker Studio ダッシュボードの表示速度も 10秒 → 1秒未満に改善される。
  5. dbt Mesh cross-project ref の安定化: access: public + latest_version: 1 で公開 API としてバージョン管理。on_schema_change: fail を上流モデルに設定することで、スキーマ変更時に下流 cross-project ref が断裂する前にビルドを失敗させて検知できる。下流は {{ ref('my_project', 'mart_member_channel_monthly', v=1) }} で参照。

実務への応用

ECサイト MOps チームの 会員ランク別キャンペーン ROAS 分析 では本パターンが直接適用できる:

  • 週次キャンペーン MTG: Looker Studio が mv_member_channel_monthly を参照 → フルスキャンなし、ダッシュボード表示 <1 秒でマーケターがリアルタイム意思決定
  • Argo Workflows: dbt source freshnessdbt rundbt test の3ステップを CronWorkflow で毎朝 JST 06:00 実行。テスト失敗時は Argo が Workflow を Failed 状態にして DataDog アラート
  • DataDog 鮮度監視: dbt source freshness --output json の結果を statsd 経由で DataDog に送信。elapsed_hours > 2 → warn(P3)、> 6 → error(P2 Slack アラート)
  • GCP コスト配賦: mart_member_channel_monthlylabels: {team: "mops", domain: "campaign"} を付与し、INFORMATION_SCHEMA.JOBS でチーム別コスト可視化(月次 FinOps レビューに活用)
  • CI/CD: Cloud Build で PR 時に dbt test --select mart_member_channel_monthly --vars '{"lookback_days": 7}' を実行。Unit Tests 失敗でマージブロック

今日のまとめ

Python ループ × BigQuery API の N+1 問題は BigQuery JOIN で 1 クエリ化し、dbt Unit Tests(1.8+)で rank_weight CASE 式・ゼロ除算ガードを自動検証、BigQuery Materialized View で下流スキャンコストを月額 $200 → $10 に削減するのが王道パターン。var('lookback_days', 90) + insert_overwrite で遅延データも自動補正される。

自己評価