C データエンジニアリング — dbt Snapshots SCD Type 2(invalidate_hard_deletes: true)× JSON_VALUE/JSON_QUERY 正規化 × BigQuery Editions INFORMATION_SCHEMA.JOBS × RESERVATIONS × ASSIGNMENTS スロット監視 × current_snapshot マクロ(MOps 会員属性変更履歴 Bad→Good 7点)

2026-07-01 (Day 86) 水曜 C: データエンジニアリング ★★★★☆ BigQuery / dbt Core 1.8+ / Snapshots / SCD Type 2 JSON_VALUE / JSON_QUERY / Editions / INFORMATION_SCHEMA

概要

📸

dbt Snapshots(SCD Type 2)で変更履歴を自動管理

全件 DELETE + INSERT で会員属性を「上書き」するアンチパターンは過去データを消滅させる。dbt Snapshots は strategy: timestamp + updated_at: changed_at で差分行のみ追記し、valid_from / valid_to を自動付与する。invalidate_hard_deletes: true で退会(物理削除)も valid_to に自動記録される。

🔍

JSON_VALUE / JSON_QUERY で payload を構造化列に変換

JSON 文字列として保存されたプロフィール変更ペイロードを SELECT * で Python に取り込み解析するパターンはスキャン無駄 + N+1 問題を生む。JSON_VALUE(payload, '$.rank') でスカラー値、JSON_QUERY(payload, '$.preferences') でネストオブジェクトを BigQuery SQL 内で直接抽出することで、必要カラムのみのスキャンが実現する。

📊

INFORMATION_SCHEMA で Editions スロット使用率を可視化

BigQuery Editions(Enterprise Plus)移行後のスロット使用量は INFORMATION_SCHEMA.JOBS × RESERVATIONS × ASSIGNMENTS を結合することで可視化できる。avg_slots_used / slot_capacity × 100 がスロット使用率で、80% 超えは待ちキュー発生リスクのサインとして DataDog アラートを設定する。

🔧

current_snapshot マクロで valid_to IS NULL を共通化

SCD Type 2 テーブルを参照する下流モデルが全て WHERE valid_to IS NULL を個別記述すると、カラム名変更時に全修正が必要になる。dbt マクロ current_snapshot('snap_member_profile') で「最新レコードの取得方法」を単一箇所に集約し、保守性を高める。

問題

ECサイトの MOps チームでは、BigQuery × dbt Core 1.8+ で 会員属性変更履歴(SCD Type 2) を管理している。以下の Python スクリプト・dbt モデルには 7つの設計上の問題 が潜んでいる。問題点を全て洗い出し、dbt Snapshots(SCD Type 2)・BigQuery JSON_VALUE/JSON_QUERY・BigQuery Editions スロット予約・INFORMATION_SCHEMA.JOBS × Reservations × Assignment・dbt Tests(generic + singular) を活用した Bad→Good リファクタリングを行え。

制約・前提条件

  • dbt Core 1.8+(Snapshots 設定ベースの YAML 形式、snapshot_meta_column_names 対応)
  • BigQuery Standard SQL(GA 版機能のみ)
  • テーブル: events.member_profile_changes(partition: changed_at TIMESTAMP, cluster: member_id)— JSON payload 列にプロフィール変更内容が格納
  • Argo Workflows CronWorkflow で毎朝 JST 06:00 実行
  • BigQuery Editions(Enterprise Plus)スロット予約環境(オンデマンド請求を廃止済み)
期待する回答形式: 問題点の列挙(番号付き)+ 改善後 dbt Snapshot YAML + dbt Snapshot SQL + dbt Tests YAML + BigQuery SQL(JSON_VALUE/JSON_QUERY 正規化)+ INFORMATION_SCHEMA スロット監視クエリ + 設計意図の説明

悪いコード (Before)

このコード・モデルには 7つのデータエンジニアリング設計の問題 が隠れています。
bad_member_profile.py + bad snapshot — 全件 DELETE+INSERT / SELECT * / JSON ハードコード / テストなし
from google.cloud import bigquery
import json

client = bigquery.Client()

def sync_member_profiles():
    # 問題①: SELECT * → JSON payload 含む全カラムスキャン(無駄コスト)
    # 問題①: Python ループで BigQuery API を複数回呼び出し(N+1)
    query = """
        SELECT *
        FROM `my-project.events.member_profile_changes`
        WHERE changed_at >= TIMESTAMP_SUB(CURRENT_TIMESTAMP(), INTERVAL 25 HOUR)
    """
    # 問題③: マジックナンバー 25 がハードコード
    rows = client.query(query).result()

    for row in rows:
        # 問題③: JSON キー名が Python 側にハードコード(SQL と二重管理)
        payload = json.loads(row["payload"])
        rank    = payload.get("rank")
        email   = payload.get("email")

        # 問題②: 全件 DELETE + INSERT(変更差分追跡不能・過去データ消滅)
        delete_q = f"""
            DELETE FROM `my-project.snapshots.member_profile_snapshot`
            WHERE member_id = '{row['member_id']}'
        """
        insert_q = f"""
            INSERT INTO `my-project.snapshots.member_profile_snapshot`
            VALUES ('{row['member_id']}', '{rank}', '{email}',
                    CURRENT_TIMESTAMP(), NULL)
        """
        client.query(delete_q).result()
        client.query(insert_q).result()
        # 問題④: スロット使用量の可視化なし(Editions 移行後のコスト効果不明)

# ── dbt 側 ──
# snapshots/snap_member_profile.sql(悪い例)
# {% snapshot snap_member_profile %}
# {{
#   config(
#     target_schema='snapshots',
#     unique_key='member_id',
#     strategy='timestamp',
#     updated_at='changed_at',
#     # 問題⑤: invalidate_hard_deletes が未設定(デフォルト false)
#     # → 退会(物理削除)が検知されず dbt_valid_to が NULL のまま
#   )
# }}
# SELECT * FROM {{ source('events', 'member_profile_changes') }}
# {% endsnapshot %}

# 問題⑥: dbt Tests 未設定
# 問題⑦: 下流モデルが全て個別に dbt_valid_to IS NULL を記述
# → mart_member_current.sql: WHERE dbt_valid_to IS NULL
# → mart_campaign_target.sql: WHERE dbt_valid_to IS NULL
# → mart_retention.sql: WHERE dbt_valid_to IS NULL  ← カラム名変更で全壊
問題点サマリー(7点)
1SELECT * + Python ループ N+1 — JSON payload 含む全カラムスキャン + BigQuery API を行ごとに呼び出し。JSON_VALUE で必要フィールドのみ抽出し BigQuery SQL 内で処理
2全件 DELETE + INSERT でスナップショット管理 — 変更差分が追跡不能・過去データが消滅。dbt Snapshots(SCD Type 2)で valid_from / valid_to を自動管理
3JSON キー名・lookback 値のハードコード散在 — Python と SQL に二重管理。JSON_VALUE(payload, '$.rank') で SQL に統一 + var('lookback_hours', 25) で一元管理
4BigQuery スロット使用量の可視化なし — Editions 移行後のコスト効果・キュー待ちリスクが不可視。INFORMATION_SCHEMA.JOBS × RESERVATIONS × ASSIGNMENTS で日次監視
5invalidate_hard_deletes: false(デフォルト) — 退会(物理削除)を検知せず valid_to が NULL のまま。true に設定し退会日を自動クローズ
6dbt Tests 未設定member_rank の想定外値・member_id 重複を検知できない。not_null, accepted_values, unique で品質ゲートを整備
7dbt_valid_to IS NULL フィルタの散在 — 下流全モデルが個別記述。カラム名変更で全壊。current_snapshot() マクロで共通化

ヒント(段階的開示)

ヒント1 — 方向性
現状コードの最大問題は「全件 DELETE + INSERT で会員属性を上書き」している点。BigQuery × dbt の正規パターンは dbt Snapshots(SCD Type 2)で valid_from / valid_to を自動管理し、変更差分のみ追記すること。また、JSON 文字列の payload を Python で解析するのではなく、BigQuery SQL 内で JSON_VALUE() を使って構造化列に変換することで、パーティション剪定とクラスタリングを活用できる。
ヒント2 — アプローチ
  • 問題①: SELECT * + Python ループ → JSON_VALUE(payload, '$.rank') で BigQuery SQL 内に処理を移動
  • 問題②: 全件 DELETE + INSERT → dbt Snapshots(strategy: timestamp, updated_at: changed_at)で差分追記
  • 問題③: マジックナンバー 25・JSON キー名散在 → {{ var('lookback_hours', 25) }} + SQL 内 JSON_VALUE で一元管理
  • 問題④: スロット可視化なし → INFORMATION_SCHEMA.JOBS × RESERVATIONS × ASSIGNMENTS で日次監視 SQL
  • 問題⑤: invalidate_hard_deletes 未設定 → invalidate_hard_deletes: true(退会を自動クローズ)
  • 問題⑥: dbt Tests なし → not_null, accepted_values, unique, relationships
  • 問題⑦: valid_to IS NULL 散在 → current_snapshot() マクロで共通化
ヒント3 — コードの骨格
-- BigQuery JSON_VALUE / JSON_QUERY の骨格(修正①③)
WITH extracted AS (
  SELECT
    member_id,
    JSON_VALUE(payload, '$.rank')         AS member_rank,   -- スカラー値
    JSON_VALUE(payload, '$.email')        AS email,
    JSON_QUERY(payload, '$.preferences')  AS preferences_json,  -- ネストオブジェクト
    changed_at
  FROM {{ source('events', 'member_profile_changes') }}
  WHERE changed_at >= TIMESTAMP_SUB(
    CURRENT_TIMESTAMP(),
    INTERVAL {{ var('lookback_hours', 25) }} HOUR  -- 修正③
  )
)
SELECT * FROM extracted

dbt Snapshots YAML 骨格(dbt Core 1.8+):

snapshots:
  - name: snap_member_profile
    relation: source('events', 'member_profile_changes')
    config:
      schema: snapshots
      unique_key: member_id
      strategy: timestamp
      updated_at: changed_at
      invalidate_hard_deletes: true   # 修正⑤: 退会を valid_to に自動記録
      snapshot_meta_column_names:
        dbt_valid_from: valid_from
        dbt_valid_to:   valid_to

INFORMATION_SCHEMA スロット監視の骨格(修正④):

SELECT
  DATE(creation_time, 'Asia/Tokyo')                AS job_date,
  reservation_id,
  SUM(total_slot_ms) / 1000 / 86400               AS avg_slots_used,
  -- 予約スロット容量と結合して使用率を計算
  ROUND(avg_slots_used / slot_capacity * 100, 1)  AS slot_utilization_pct
FROM `region-asia-northeast1`.INFORMATION_SCHEMA.JOBS
WHERE DATE(creation_time, 'Asia/Tokyo') >= DATE_SUB(CURRENT_DATE('Asia/Tokyo'), INTERVAL 7 DAY)
  AND state = 'DONE' AND error_result IS NULL
GROUP BY job_date, reservation_id

問題点分析(7点)

#問題点分類改善方法
1SELECT * + Python ループ N+1(JSON 全スキャン + 行ごとに BQ API)パフォーマンスJSON_VALUE で必要フィールドのみ BigQuery SQL 内で抽出
2全件 DELETE + INSERT(変更差分追跡不能・過去データ消滅)データ品質dbt Snapshots SCD Type 2(valid_from / valid_to 自動管理)
3JSON キー名・lookback 値が Python と SQL に散在保守性var('lookback_hours', 25) + SQL 内 JSON_VALUE で一元管理
4BigQuery Editions スロット使用量の可視化なしコスト可視性INFORMATION_SCHEMA.JOBS × RESERVATIONS × ASSIGNMENTS で日次監視
5invalidate_hard_deletes: false(退会検知なし)データ品質invalidate_hard_deletes: true(退会日を valid_to に自動クローズ)
6dbt Tests 未設定(想定外ランク・重複 member_id を検知できない)品質not_null, accepted_values, unique, relationships
7valid_to IS NULL フィルタを下流全モデルが個別記述保守性current_snapshot() マクロで共通化

設計図 — Bad vs Good データフロー(SCD Type 2 × JSON 正規化)

Bad(変更前)— 全件 DELETE+INSERT / SELECT * / テストなし events.member_profile_changes(JSON payload) ① SELECT * → JSON payload 含む全カラムスキャン(無駄コスト) ⚠️ Python for ループ × BigQuery API → 行ごとに DELETE+INSERT(N+1) ③ lookback 25H・JSON キー名が Python 側にハードコード(SQL との二重管理) Python(Argo Workflows Pod)— 行ごとに BQ API 呼び出し ① json.loads(row["payload"]) → Python 側で JSON 解析(BQ の強みを活かせない) ② DELETE + INSERT を 1 行ずつ実行 → Editions スロット消費 × 処理時間 数分 ④ スロット使用量の可視化なし(INFORMATION_SCHEMA 未活用) ⑤ invalidate_hard_deletes 未設定 → 退会会員が valid_to=NULL のまま残る ⑥ dbt Tests なし → member_rank に想定外値が混入しても検知できない member_profile_snapshot(全件上書き) ② DELETE + INSERT で上書き → SCD Type 1(履歴ゼロ) ⚠️ キャンペーン配信時点の会員ランクを遡及参照できない ⚠️ 退会済み会員が dbt_valid_to=NULL のまま「現役」として扱われ続ける 下流 mart / analytics モデル(テストなし・フィルタ散在) ⑦ mart_member_current.sql: WHERE dbt_valid_to IS NULL ⑦ mart_campaign_target.sql: WHERE dbt_valid_to IS NULL ⚠️ カラム名を valid_to に変更すると全下流モデルが一斉に壊れる ⑥ 想定外ランク値("vip_unknown")が混入してもビルド継続 Good(変更後)— dbt Snapshots SCD Type 2 × JSON_VALUE × IS events.member_profile_changes({{ source() }}) 修正①③: JSON_VALUE(payload, '$.rank') で必要5カラムのみ抽出(スキャン大幅削減) 修正③: WHERE changed_at >= TIMESTAMP_SUB(..., INTERVAL {{ var('lookback_hours', 25) }} HOUR) JSON_QUERY(payload, '$.preferences') でネストオブジェクトを JSON 文字列で保持 partition_by: changed_at → パーティション剪定 / cluster_by: member_id → 高速ルックアップ dbt Snapshots(SCD Type 2)— snap_member_profile 修正②: strategy: timestamp + updated_at: changed_at → 差分行のみ追記 修正②: valid_from / valid_to(snapshot_meta_column_names でエイリアス)自動管理 修正⑤: invalidate_hard_deletes: true → 退会(物理削除)を valid_to に自動クローズ partition_by: valid_from(timestamp/day)× cluster_by: [member_id, member_rank] on_schema_change: fail → スキーマ変更を即検知してビルドを失敗させる dbt snapshot → dbt test --select snap_member_profile を Argo Workflows で直列実行 dbt Tests(修正⑥) not_null: member_id, valid_from accepted_values: member_rank (bronze/silver/gold/platinum) unique: member_id(最新レコード) CI: dbt test --select snap_member_profile → マージ前に自動検証 IS スロット監視(修正④) INFORMATION_SCHEMA.JOBS × RESERVATIONS × ASSIGNMENTS slot_utilization_pct 算出 80% 超え → DataDog P3 alert 100% 超え → P2 Slack alert Argo CronWorkflow 日次実行 current_snapshot() マクロ(修正⑦) 修正⑦: current_snapshot('snap_member_profile') → WHERE valid_to IS NULL を共通化 カラム名変更(dbt_valid_to → valid_to)はマクロ 1 箇所のみ修正 退会会員(valid_to IS NOT NULL)は自動的に除外 下流: mart_member_current / mart_campaign_target / mart_retention が共通マクロを呼ぶだけ 下流 mart / analytics(SCD Type 2 ポイントインタイム参照) 配信時点ランク: WHERE valid_from <= sent_at AND (valid_to > sent_at OR valid_to IS NULL) dbt Mesh cross-project ref: access: public + latest_version: 1 で analytics PJ へ公開 DataDog: dbt source freshness → warn 2h / error 6h + スロット使用率 80% アラート 修正

模範解答

snapshots/snap_member_profile.yml — dbt Snapshots SCD Type 2(修正②⑤)

# snapshots/snap_member_profile.yml
# dbt Core 1.8+ — 設定ベース YAML 形式({% snapshot %} ブロック廃止)
snapshots:
  - name: snap_member_profile
    description: "会員プロフィール変更履歴 SCD Type 2(rank/email/prefecture/preferences)"
    relation: source('events', 'member_profile_changes')
    config:
      schema: snapshots
      database: my-gcp-project
      # 修正②: SCD Type 2 — member_id をキーに変更差分のみ追記
      unique_key: member_id
      strategy: timestamp
      updated_at: changed_at          # changed_at の変化で差分を検知
      # 修正⑤: 物理削除(退会)を valid_to に自動クローズ
      invalidate_hard_deletes: true
      # カラム名を読みやすいエイリアスに変換(dbt_valid_from → valid_from)
      snapshot_meta_column_names:
        dbt_valid_from: valid_from
        dbt_valid_to:   valid_to
        dbt_scd_id:     scd_id
        dbt_updated_at: updated_at_snapshot
      # BigQuery パーティション + クラスタリング設定
      partition_by:
        field: valid_from
        data_type: timestamp
        granularity: day              # 日次パーティション(valid_from でフィルタ可能)
      cluster_by: [member_id, member_rank]  # ルックアップ高速化
      on_schema_change: fail          # スキーマ変更を即検知

    columns:
      - name: member_id
        description: "会員 ID(一意キー)"
        tests:
          - not_null
      - name: member_rank
        description: "会員ランク(bronze / silver / gold / platinum)"
        tests:
          - not_null
          - accepted_values:
              values: ["bronze", "silver", "gold", "platinum"]
              # 修正⑥: 想定外ランクが来た場合にビルドを失敗させる
      - name: email
        tests: [not_null]
      - name: valid_from
        tests: [not_null]
      # valid_to: 最新レコードは NULL(退会済みは NOT NULL)— not_null テスト不可

snapshots/snap_member_profile.sql — JSON_VALUE/JSON_QUERY 正規化(修正①③)

-- snapshots/snap_member_profile.sql
-- dbt Core 1.8+ | BigQuery Standard SQL
-- 会員プロフィール変更イベント(JSON payload)を SCD Type 2 で管理
-- ※ dbt 1.8+ では YAML 設定ベース。このファイルは SELECT のみ記述

-- 修正①③: SELECT * → 必要フィールドのみ JSON_VALUE/JSON_QUERY で抽出
-- 修正③: JSON キー名を SQL 内に統一(Python 側のハードコード廃止)
WITH extracted AS (
  SELECT
    member_id,

    -- ── スカラー値: JSON_VALUE(文字列を返す)────────────────
    JSON_VALUE(payload, '$.rank')             AS member_rank,   -- bronze/silver/gold/platinum
    JSON_VALUE(payload, '$.email')            AS email,
    JSON_VALUE(payload, '$.prefecture')       AS prefecture,
    JSON_VALUE(payload, '$.age_group')        AS age_group,     -- 20s/30s/40s/50s_plus

    -- ── ネストオブジェクト: JSON_QUERY(JSON 文字列を返す)──
    -- preferences: {"newsletter": true, "sms": false, "push": true}
    JSON_QUERY(payload, '$.preferences')      AS preferences_json,

    -- ── 数値スカラー: JSON_VALUE → SAFE_CAST で型変換 ───────
    SAFE_CAST(
      JSON_VALUE(payload, '$.lifetime_orders') AS INT64
    )                                         AS lifetime_orders,

    changed_at                                -- SCD Type 2 の updated_at キー

  FROM {{ source('events', 'member_profile_changes') }}
  -- 修正③: var() で lookback 管理(本番: 25h で JST 日跨ぎ + Pub/Sub 遅延 1h 対応)
  WHERE changed_at >= TIMESTAMP_SUB(
    CURRENT_TIMESTAMP(),
    INTERVAL {{ var('lookback_hours', 25) }} HOUR
  )
  -- 重複イベント対策: 同一 member_id × changed_at の最新を使用
  QUALIFY ROW_NUMBER() OVER (
    PARTITION BY member_id, DATE(changed_at, 'Asia/Tokyo')
    ORDER BY changed_at DESC
  ) = 1
)

SELECT
  member_id,
  member_rank,
  email,
  prefecture,
  age_group,
  preferences_json,
  lifetime_orders,
  changed_at
FROM extracted

models/snapshots/snap_member_profile.yml — dbt Tests(修正⑥)

# models/snapshots/snap_member_profile.yml(別途 schema tests)
# dbt Core 1.8+ — generic tests + singular tests
version: 2

models:
  # ステージングモデル(JSON 正規化済み・最新レコードのみ)
  - name: stg_member_profile_current
    description: "会員プロフィール現在状態(snap_member_profile の valid_to IS NULL ビュー)"
    columns:
      - name: member_id
        tests:
          - not_null
          - unique   # 現在状態: member_id は一意のはず

      - name: member_rank
        tests:
          - not_null
          - accepted_values:
              values: ["bronze", "silver", "gold", "platinum"]
              # 想定外値が来るとビルドが失敗し CI でブロック

      - name: email
        tests:
          - not_null

      - name: valid_from
        tests: [not_null]

# ── Singular Tests(カスタム SQL テスト)────────────────────
# tests/assert_no_overlapping_snapshot_periods.sql
# member_id ごとに valid_from と valid_to が重複していないことを確認
#
# SELECT member_id
# FROM (
#   SELECT
#     a.member_id,
#     a.valid_from AS a_from, a.valid_to AS a_to,
#     b.valid_from AS b_from, b.valid_to AS b_to
#   FROM {{ ref('snap_member_profile') }} a
#   JOIN {{ ref('snap_member_profile') }} b
#     ON a.member_id = b.member_id
#     AND a.scd_id != b.scd_id
#     AND a.valid_from < COALESCE(b.valid_to, TIMESTAMP('9999-12-31'))
#     AND COALESCE(a.valid_to, TIMESTAMP('9999-12-31')) > b.valid_from
# )
# WHERE member_id IS NOT NULL
# -- レコードが 0 件なら PASS / 1 件以上なら FAIL

# ── Argo Workflows での実行順序 ──────────────────────────────
# 1. dbt source freshness --select source:events.member_profile_changes
# 2. dbt snapshot --select snap_member_profile
# 3. dbt test --select snap_member_profile
# 4. dbt run --select stg_member_profile_current+
# 5. dbt test --select stg_member_profile_current

BigQuery INFORMATION_SCHEMA スロット使用率監視(修正④)

-- BigQuery Editions スロット使用量 × Reservation 割り当て日次監視
-- 実行: Argo Workflows 日次 CronWorkflow(JST 06:10 = dbt 実行後)
-- 出力: DataDog statsd 経由でメトリクス送信 → slot_utilization_pct > 80% で P3 alert

WITH
job_slots AS (
  SELECT
    DATE(creation_time, 'Asia/Tokyo')                       AS job_date,
    reservation_id,
    job_type,
    -- 平均使用スロット数(total_slot_ms = ジョブ全体のスロット×ミリ秒の合計)
    SUM(total_slot_ms) / 1000                               AS total_slot_seconds,
    -- 1 日(86400 秒)に対する平均使用スロット数
    ROUND(SUM(total_slot_ms) / 1000.0 / 86400.0, 1)        AS avg_slots_used,
    COUNT(*)                                                AS job_count,
    ROUND(SUM(total_bytes_processed) / POW(1024.0, 4), 4)  AS total_tb_processed,
    -- スロット消費上位ジョブを確認(デバッグ用)
    ARRAY_AGG(
      STRUCT(job_id, total_slot_ms, query)
      ORDER BY total_slot_ms DESC LIMIT 3
    )                                                       AS top_3_jobs
  FROM `region-asia-northeast1`.INFORMATION_SCHEMA.JOBS
  WHERE
    -- 直近 7 日間(スロット使用トレンド確認)
    DATE(creation_time, 'Asia/Tokyo')
      >= DATE_SUB(CURRENT_DATE('Asia/Tokyo'), INTERVAL 7 DAY)
    AND state = 'DONE'
    AND error_result IS NULL          -- エラージョブを除外(スロット消費が不正確)
    AND job_type = 'QUERY'            -- クエリジョブのみ(DDL/LOAD 除外)
  GROUP BY job_date, reservation_id, job_type
),

reservation_capacity AS (
  -- 予約スロット容量(購入済みスロット数)
  SELECT
    reservation_name,
    slot_capacity,
    edition,                          -- ENTERPRISE_PLUS / ENTERPRISE / STANDARD
    ignore_idle_slots                 -- true: 他予約のアイドルスロットを借用可能
  FROM `region-asia-northeast1`.INFORMATION_SCHEMA.RESERVATIONS
),

assignment AS (
  -- プロジェクト ↔ 予約の割り当て(プロジェクト別コスト配賦に使用)
  SELECT
    reservation_name,
    assignee_id,                      -- プロジェクト ID または folder ID
    job_type AS assigned_job_type
  FROM `region-asia-northeast1`.INFORMATION_SCHEMA.ASSIGNMENTS
  WHERE assignee_type = 'PROJECT'
)

SELECT
  j.job_date,
  j.reservation_id,
  r.edition,
  r.slot_capacity,
  j.avg_slots_used,
  -- スロット使用率(100% 超え = 待ちキューが発生してクエリが遅延)
  ROUND(j.avg_slots_used / NULLIF(r.slot_capacity, 0) * 100.0, 1) AS slot_utilization_pct,
  j.job_count,
  ROUND(j.total_tb_processed, 4)                                   AS total_tb_processed,
  j.top_3_jobs,
  a.assignee_id                                                    AS project_id,
  -- アラート判定(DataDog カスタムメトリクスで閾値監視)
  CASE
    WHEN j.avg_slots_used / NULLIF(r.slot_capacity, 0) >= 1.0 THEN 'CRITICAL'  -- 100%: P2
    WHEN j.avg_slots_used / NULLIF(r.slot_capacity, 0) >= 0.8 THEN 'WARNING'   -- 80%: P3
    ELSE 'OK'
  END AS alert_level
FROM job_slots AS j
LEFT JOIN reservation_capacity AS r
  ON j.reservation_id LIKE CONCAT('%', r.reservation_name)
LEFT JOIN assignment AS a
  ON r.reservation_name = a.reservation_name
ORDER BY j.job_date DESC, slot_utilization_pct DESC

macros/current_snapshot.sql — valid_to IS NULL の共通化(修正⑦)

-- macros/current_snapshot.sql
-- 修正⑦: dbt_valid_to IS NULL フィルタを共通化
-- 使用例: SELECT * FROM {{ current_snapshot('snap_member_profile') }}

{% macro current_snapshot(snapshot_name, as_of=None) %}
  {%- if as_of -%}
    {# ポイントインタイムクエリ: キャンペーン配信時点の会員ランクを取得 #}
    SELECT *
    FROM {{ ref(snapshot_name) }}
    WHERE valid_from <= TIMESTAMP('{{ as_of }}')
      AND (
        valid_to > TIMESTAMP('{{ as_of }}')
        OR valid_to IS NULL
      )
  {%- else -%}
    {# 最新状態のみ(退会会員は valid_to NOT NULL なので自動除外)#}
    SELECT *
    FROM {{ ref(snapshot_name) }}
    WHERE valid_to IS NULL
  {%- endif -%}
{% endmacro %}


-- ── 使用例 1: 現在の会員状態 ──────────────────────────────────
-- models/mart/mart_member_current.sql
-- SELECT
--   member_id,
--   member_rank,
--   email,
--   valid_from AS rank_changed_at
-- FROM {{ current_snapshot('snap_member_profile') }}


-- ── 使用例 2: キャンペーン配信時点のランク(ポイントインタイム)──
-- models/mart/mart_campaign_attributed_revenue.sql
-- WITH campaign_sends AS (
--   SELECT member_id, campaign_id, sent_at FROM {{ ref('stg_campaign_sends') }}
-- ),
-- member_at_send AS (
--   -- 配信時点の会員ランクを取得
--   SELECT * FROM {{ current_snapshot('snap_member_profile', as_of='2026-06-01') }}
-- )
-- SELECT
--   s.campaign_id,
--   m.member_rank,
--   COUNT(DISTINCT s.member_id) AS sent_count
-- FROM campaign_sends AS s
-- JOIN member_at_send AS m USING (member_id)
-- GROUP BY s.campaign_id, m.member_rank

dbt_project.yml — vars + Snapshots 設定(修正③)

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

# 修正③: マジックナンバー一元管理
vars:
  # lookback_hours: JST 日跨ぎ(1h)+ Pub/Sub 遅延(1h)+ バッファ(23h)= 25h
  lookback_hours: 25    # CI 実行: dbt run --vars '{"lookback_hours": 2}'

# Snapshots グローバル設定
snapshots:
  my_project:
    +target_schema: snapshots
    +invalidate_hard_deletes: true    # 修正⑤: 全 Snapshot で物理削除を検知

models:
  my_project:
    staging:
      +materialized: view
    mart:
      +materialized: table
      +partition_by:
        field: valid_from
        data_type: timestamp
        granularity: day

# ── Argo Workflows CronWorkflow(JST 06:00 実行)──────────────
# schedule: "0 21 * * *"   # UTC 21:00 = JST 06:00
# steps:
#   1. dbt source freshness --select source:events.member_profile_changes
#   2. dbt snapshot --select snap_member_profile
#   3. dbt test --select snap_member_profile
#   4. dbt run --select stg_member_profile_current+
#   5. dbt test --select stg_member_profile_current+
#   6. [slot_monitor.sql を BQ に投げて結果を DataDog statsd 送信]

ポイント解説

  1. JSON_VALUE vs JSON_QUERY の使い分け: JSON_VALUE(payload, '$.rank') はスカラー値(文字列・数値・ブール)を STRING 型で返す。ネストしたオブジェクトや配列は JSON_QUERY(payload, '$.preferences') で JSON 文字列のまま取得。BigQuery の JSON 型カラムなら payload.rank の記法も使えるが、STRING 型 JSON の場合は JSON_VALUE / JSON_QUERY が必須。数値変換は SAFE_CAST(JSON_VALUE(payload, '$.lifetime_orders') AS INT64) でキャスト。
  2. dbt Snapshots SCD Type 2 の動作原理: strategy: timestamp + updated_at: changed_at の設定で、changed_at が前回より新しい行を「変更あり」と判定し、旧バージョンの valid_to を現在時刻でクローズして新バージョンを追記する。invalidate_hard_deletes: true を設定すると、ソーステーブルから行が消えた際(会員退会など)に valid_to を現在時刻で自動クローズする。
  3. INFORMATION_SCHEMA.JOBS × RESERVATIONS × ASSIGNMENTS の結合: BigQuery Editions 環境では、JOBS ビューの reservation_id(形式: project:region.reservation_name)を RESERVATIONSreservation_name と LIKE 結合することでスロット使用率を計算できる。total_slot_ms / 1000 / 86400 で日平均使用スロット数を算出。使用率が 80% 超えはスロット購入見直しの目安。
  4. snapshot_meta_column_names によるカラム名エイリアス: デフォルトの dbt_valid_from / dbt_valid_to より valid_from / valid_to の方が下流モデルで読みやすい。dbt Core 1.8+ の snapshot_meta_column_names 設定でエイリアスを指定すると、dbt が自動生成するメタカラム名をカスタマイズできる。これにより current_snapshot() マクロの WHERE valid_to IS NULL も一貫して使える。
  5. QUALIFY + ROW_NUMBER() で重複イベントを除去: Pub/Sub の at-least-once 配信や BQ Streaming Insert のリトライにより、同一 member_id × 同日の変更イベントが複数到着することがある。QUALIFY ROW_NUMBER() OVER (PARTITION BY member_id, DATE(changed_at, 'Asia/Tokyo') ORDER BY changed_at DESC) = 1 で最新イベントのみを抽出し、dbt Snapshot の重複追記を防ぐ。

実務への応用

ECサイト MOps チームの 会員セグメント変更追跡 × キャンペーン ROAS 分析 では本パターンが直接適用できる:

  • キャンペーン配信時点のランク取得: current_snapshot('snap_member_profile', as_of='2026-06-01') でポイントインタイムクエリ。「配信時点でゴールドだった会員の購買率」を正確に算出できる
  • 退会会員の自動除外: invalidate_hard_deletes: true により退会処理(ソース行削除)が翌朝 CronWorkflow 実行時に valid_to で自動クローズ → current_snapshot()valid_to IS NULL フィルタで自然に除外
  • BigQuery Editions コスト最適化: slot_utilization_pct を日次で DataDog に送信。70% 以下が続く場合はスロット削減提案、90% 超えが続く場合はスロット追加を CFO に提案
  • ランク変動分析: LAG(member_rank) OVER (PARTITION BY member_id ORDER BY valid_from) で前ランクを取得し、ブロンズ→ゴールドへの昇格イベントをキャンペーントリガーに活用
  • dbt Mesh 公開: snap_member_profileaccess: public + latest_version: 1 で analytics プロジェクトへ公開し、BI チームが直接参照できる安定 API として提供

今日のまとめ

JSON payload の JSON_VALUE/JSON_QUERY 正規化 + dbt Snapshots SCD Type 2(invalidate_hard_deletes: true)の組み合わせが会員属性変更履歴の王道パターン。BigQuery Editions 環境では INFORMATION_SCHEMA.JOBS × RESERVATIONS × ASSIGNMENTS でスロット使用率を日次監視し、current_snapshot() マクロで valid_to IS NULL フィルタを共通化することで下流モデルの保守性を高める。

自己評価