C データエンジニアリング — BigQuery ROW_NUMBER ラスト・タッチ attribution × dbt source freshness × PIVOT 演算子 × SELECT * スキャン削減 × N+1 Python ループ廃止(MOps キャンペーン attribution Bad→Good)

2026-06-17 (Day 72) 水曜 C: データエンジニアリング ★★★★☆ BigQuery / dbt Core 1.8+ / ウィンドウ関数 / PIVOT Attribution Modeling / dbt Freshness / dbt Mesh

概要

📊

ROW_NUMBER() ウィンドウ関数でラスト・タッチ attribution

Python ループで行ごとに attribution を判定するアンチパターンは N+1 BigQuery API 呼び出しを引き起こす。ROW_NUMBER() OVER(PARTITION BY order_id ORDER BY click_time DESC) で同一 order に紐づく複数クリックのうち最後のクリックのみを採用し、BigQuery の分散処理エンジンに全処理を委譲する。処理時間が数時間から数十秒に短縮される。

🌿

dbt source() + freshness で依存グラフと鮮度監視を統合

テーブル名のハードコードは dbt の依存グラフを壊し、テーブル移行時に全モデルを grep 置換する羽目になる。{{ source('events', 'email_clicks') }} で依存を明示し、loaded_at_field: _ingested_at + warn_after: {count: 2, period: hour} で Pub/Sub → BQ パイプライン遅延を自動検知する。

💰

SELECT * 廃止でスキャン85%削減

BigQuery は列指向ストレージのため SELECT * は全カラムをスキャンする。50カラムのテーブルで必要6カラムのみ取得すると月額 $150 → $22.5(85%削減)。パーティション剪定(WHERE event_date >= ...)と組み合わせると更に大幅削減が可能。

🔄

BigQuery PIVOT 演算子で pandas OOM を解消

500キャンペーン × 月次データを pandas で pivot_table するとメモリ上に全データを展開するため Argo Workflows Pod(512Mi)がメモリ不足で失敗する。BigQuery の PIVOT 演算子は BigQuery クラスター上で処理するためメモリ消費がゼロで、Pod に影響を与えない。

問題

ECサイトの MOps チームでは、BigQuery 上でメールキャンペーンの売上貢献分析を行っている。現在の SQL・dbt モデルには 7つの設計上の問題 が潜んでいる。問題点を全て洗い出し、BigQuery ウィンドウ関数・PIVOT・dbt Mesh cross-project ref・dbt Freshness・INFORMATION_SCHEMA.PARTITIONS を活用した Bad→Good リファクタリングを行え。

制約・前提条件

  • BigQuery Standard SQL(GA版機能のみ使用可)
  • dbt Core 1.8+(Unit Tests / Mesh cross-project ref 対応)
  • テーブル: events.email_clicks(partition: event_date DATE, cluster: campaign_id), orders.order_items(partition: order_date DATE
  • アトリビューション窓: クリック後 7日以内の注文を貢献とみなす(設定値として外部化すること)
  • Argo Workflows の cronworkflow で毎朝 JST 07:00 実行
  • dbt source freshness + DataDog メトリクスへの鮮度アラートも設定すること
期待する回答形式: 問題点の列挙(番号付き)+ 改善後 SQL(dbt モデル)+ sources.yml(freshness 設定)+ 設計意図の説明

悪いコード (Before)

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

client = bigquery.Client()

def run_attribution():
    # 問題①: SELECT * → 全50カラムスキャン(月額 $150 無駄)
    query = "SELECT * FROM `events.email_clicks`"
    clicks = client.query(query).to_dataframe()

    # 問題②③: Python ループ + マジックナンバー 7
    attributed = []
    for _, click in clicks.iterrows():
        # 問題②: 1行ごとに BigQuery API を呼ぶ(N+1 クエリ問題)
        orders_q = f"""
            SELECT order_id, revenue
            FROM `orders.order_items`
            WHERE user_id = '{click['user_id']}'
              AND order_date BETWEEN '{click['event_date']}'
              AND DATE_ADD('{click['event_date']}', INTERVAL 7 DAY)
        """
        # 問題③: マジックナンバー 7 が SQL と Python に散在
        orders = client.query(orders_q).to_dataframe()
        for _, order in orders.iterrows():
            attributed.append({
                "campaign_id": click["campaign_id"],
                "order_id": order["order_id"],
                "revenue": order["revenue"],
            })

    df = pd.DataFrame(attributed)

    # 問題④: ラスト・タッチ attribution の処理なし → 重複カウント
    result = df.groupby("campaign_id")["revenue"].sum()

    # 問題⑦: pandas pivot → 500キャンペーンで Pod OOM(512Mi)
    pivot = df.pivot_table(
        index="campaign_id", values="revenue", aggfunc="sum"
    )
    pivot.to_gbq("mart.campaign_attribution_pivot", if_exists="replace")

# ── dbt モデル(問題⑤⑥)──
# models/mart/mart_campaign_revenue.sql
# 問題⑤: テーブル名をハードコード(source() を使わない)
# SELECT campaign_id, SUM(revenue) AS revenue
# FROM `my-project.events.email_clicks`   ← ハードコード
# GROUP BY 1
# 問題⑥: dbt source freshness 未設定 → データ遅延を検知できない
問題点サマリー(7点)
1SELECT * によるフルカラムスキャン — 50カラム全取得で月額 $150。必要6カラムのみ取得で $22.5(85%削減)
2Python ループによる N+1 BigQuery API 問題 — 100万行クリックに対し100万回クエリ。BigQuery JOIN に委譲すれば1クエリで完結
3マジックナンバー 7 の散在 — SQL と Python の2箇所にハードコード。dbt var() で一元管理
4ラスト・タッチ attribution の未定義 — 同一 order_id に複数クリックが紐づく場合に重複カウント。ROW_NUMBER() OVER で解決
5dbt source() 未使用・テーブル名ハードコード — 依存グラフが機能せず移行時にすべてのモデルを手動修正
6dbt source freshness 未設定 — Pub/Sub → BQ パイプライン遅延(2時間超)を検知できない
7pandas PIVOT → Pod OOM — 500キャンペーン × 月次を pandas でメモリ展開。BigQuery PIVOT 演算子に移行

ヒント(段階的開示)

ヒント1 — 方向性
現状コードは Python で「行ごとに BigQuery に問い合わせる」N+1 パターンが最大の問題。BigQuery の設計思想は「大量データを一括で SQL 処理する」であり、ループは禁忌。JOIN + ウィンドウ関数で全処理を BigQuery に委譲すること。dbt では {{ source() }} でテーブル依存を管理し、sources.yml の freshness で鮮度アラートを自動化する。
ヒント2 — アプローチ
  • 問題①: SELECT * → 必要カラム6つのみ明示(BigQuery は列指向なのでカラム数でコストが変わる)
  • 問題②: Python for ループ + クエリ → BigQuery JOIN 1本で解決(user_id + 日付範囲で INNER JOIN)
  • 問題③: マジックナンバー 7{{ var('attribution_window_days', 7) }}(dbt_project.yml で管理)
  • 問題④: 重複カウント → ROW_NUMBER() OVER(PARTITION BY order_id ORDER BY click_time DESC) = 1
  • 問題⑤: ハードコード → {{ source('events', 'email_clicks') }}
  • 問題⑥: freshness なし → loaded_at_field: _ingested_at + warn_after: {count: 2, period: hour}
  • 問題⑦: pandas PIVOT → BigQuery PIVOT(SUM(revenue) FOR ym IN (...))
ヒント3 — コードの骨格
-- ウィンドウ関数でラスト・タッチ attribution(ヒント: ROW_NUMBER)
WITH ranked_clicks AS (
  SELECT
    c.campaign_id,
    c.user_id,
    c.click_time,
    o.order_id,
    o.revenue,
    -- ラスト・タッチ: 同一 order_id に複数クリックが紐づく場合は最後のクリックを採用
    ROW_NUMBER() OVER (
      PARTITION BY o.order_id
      ORDER BY c.click_time DESC
    ) AS click_rank
  FROM {{ source('events', 'email_clicks') }} AS c
  INNER JOIN {{ source('orders', 'order_items') }} AS o
    ON c.user_id = o.user_id
    AND o.order_date BETWEEN c.click_date
                         AND DATE_ADD(c.click_date, INTERVAL {{ var('attribution_window_days', 7) }} DAY)
)
SELECT campaign_id, SUM(revenue) AS attributed_revenue
FROM ranked_clicks
WHERE click_rank = 1   -- ラスト・タッチのみ集計
GROUP BY campaign_id;

問題点分析(7点)

#問題点分類改善方法
1SELECT * フルカラムスキャンコスト必要6カラムのみ明示(85%削減)
2Python N+1 BigQuery API ループパフォーマンスBigQuery JOIN 1本に集約
3マジックナンバー 7 の散在保守性dbt var('attribution_window_days', 7)
4ラスト・タッチ attribution 未定義正確性ROW_NUMBER() OVER + WHERE click_rank=1
5dbt source() 未使用・ハードコード依存管理{{ source('events', 'email_clicks') }}
6dbt source freshness 未設定可観測性loaded_at_field + warn/error_after
7pandas PIVOT → Pod OOM信頼性BigQuery PIVOT 演算子

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

Bad(変更前)— Python ループ + pandas BigQuery: events.email_clicks ① SELECT * → 50カラム全スキャン(月額 $150) ⚠️ dbt モデルでテーブル名ハードコード(source() なし) 100万行 to_dataframe() Python ループ(Argo Workflows Pod) ② for _, click in clicks.iterrows(): → N+1 クエリ(100万回) ⚠️ 1クリック → 1 BigQuery API 呼び出し → 処理時間: 数時間 ③ INTERVAL 7 DAY — マジックナンバーが SQL と Python に散在 ④ ラスト・タッチ未処理 → 複数クリックが同一 order で重複カウント pandas pivot_table(Pod メモリ上) ⑦ 500キャンペーン × 月次データをメモリに展開 ⚠️ Argo Pod 512Mi → OOMKilled で Workflow 失敗 mart.campaign_attribution_pivot(不正確) 重複 attribution / freshness 監視なし / 処理時間不安定 ⑥ データ遅延(6時間超)を検知できず、マーケが古いデータで意思決定 Good(変更後)— BigQuery SQL 一括処理 BigQuery: {{ source('events', 'email_clicks') }} 修正①: campaign_id, user_id, event_date, click_time の6カラムのみ(月額 $22.5) 修正⑤: dbt source() で依存グラフ管理 / 修正⑥: freshness warn 2h / error 6h SQL JOIN(1クエリ) BigQuery: JOIN + ROW_NUMBER() ウィンドウ関数 修正②: Python ループ廃止 → BigQuery の分散処理エンジンに委譲 INNER JOIN ON user_id + order_date BETWEEN click_date AND DATE_ADD(..., INTERVAL {{ var('attribution_window_days', 7) }} DAY) 修正③: dbt var('attribution_window_days', 7) で一元管理(dbt_project.yml) 修正④: ROW_NUMBER() OVER(PARTITION BY order_id ORDER BY click_time DESC) AS click_rank WHERE click_rank = 1 → ラスト・タッチのみ集計(重複排除)/ 処理時間: 数十秒 BigQuery PIVOT 演算子(クラスター上で処理) 修正⑦: PIVOT(SUM(attributed_revenue) FOR ym IN (...)) — Pod メモリ使用ゼロ 500キャンペーン × 6ヶ月分でも OOM なし / dbt モデルとして管理 mart.campaign_attribution(正確・低コスト・監視付き) ラスト・タッチ attribution 正確 / スキャン85%削減 / 処理時間: 数十秒 dbt freshness → DataDog アラート / dbt source() で依存グラフ可視化 マーケターが最新データ(鮮度保証)でキャンペーン ROAS を意思決定 dbt source freshness + DataDog 鮮度アラート CronWorkflow で毎時 `dbt source freshness --output json` を実行 elapsed_hours > 2h → warn / > 6h → error → DataDog metric → Slack アラート Pub/Sub → BQ streaming パイプライン障害を自動検知 修正

模範解答

mart_campaign_attribution.sql — 改善後 dbt モデル

-- models/mart/mart_campaign_attribution.sql
-- dbt Core 1.8+ | BigQuery Standard SQL
-- Attribution: ラスト・タッチ(クリック後 {{ var('attribution_window_days', 7) }} 日以内)
{{
  config(
    materialized     = 'incremental',
    partition_by     = {'field': 'click_date', 'data_type': 'date'},
    cluster_by       = ['campaign_id'],
    on_schema_change = 'fail',
  )
}}

WITH
-- 修正①: 必要カラムのみ取得(スキャン85%削減)
-- 修正⑤: {{ source() }} で依存グラフを明示
email_clicks AS (
  SELECT
    campaign_id,
    user_id,
    DATE(event_timestamp, 'Asia/Tokyo') AS click_date,   -- JST 変換
    TIMESTAMP_TRUNC(event_timestamp, SECOND) AS click_time
  FROM {{ source('events', 'email_clicks') }}             -- 修正⑤: source() 使用
  {% if is_incremental() %}
    -- インクリメンタル: 前日分 + attribution_window_days ルックバック(遅延データ対応)
    WHERE event_date >= DATE_SUB(
      DATE(CURRENT_TIMESTAMP(), 'Asia/Tokyo'),
      INTERVAL {{ var('attribution_window_days', 7) }} + 1 DAY   -- 修正③: var() で管理
    )
  {% endif %}
),

order_items AS (
  SELECT
    user_id,
    order_id,
    DATE(order_timestamp, 'Asia/Tokyo') AS order_date,
    revenue
  FROM {{ source('orders', 'order_items') }}
  {% if is_incremental() %}
    WHERE order_date >= DATE_SUB(
      DATE(CURRENT_TIMESTAMP(), 'Asia/Tokyo'),
      INTERVAL {{ var('attribution_window_days', 7) }} + 1 DAY
    )
  {% endif %}
),

-- 修正②: Python ループ廃止 → BigQuery JOIN で一括処理(N+1 問題を解消)
-- 修正③: マジックナンバー → var('attribution_window_days', 7) で一元管理
joined AS (
  SELECT
    c.campaign_id,
    c.click_date,
    c.click_time,
    o.order_id,
    o.revenue,
    -- 修正④: 同一 order_id に複数クリックが紐づく場合は最後のクリックを採用(ラスト・タッチ)
    ROW_NUMBER() OVER (
      PARTITION BY o.order_id
      ORDER BY c.click_time DESC    -- 最後のクリックを優先
    ) AS click_rank
  FROM email_clicks AS c
  INNER JOIN order_items AS o
    ON  c.user_id    = o.user_id
    -- attribution 窓: クリック日から N 日以内の注文を対象とする
    AND o.order_date BETWEEN c.click_date
                         AND DATE_ADD(
                               c.click_date,
                               INTERVAL {{ var('attribution_window_days', 7) }} DAY
                             )
),

-- 修正④: ラスト・タッチのみ残す(click_rank = 1)
last_touch AS (
  SELECT
    campaign_id,
    click_date,
    order_id,
    revenue
  FROM joined
  WHERE click_rank = 1    -- ラスト・タッチ attribution: 重複排除
),

-- キャンペーン別日次集計
campaign_daily AS (
  SELECT
    campaign_id,
    click_date,
    COUNT(DISTINCT order_id) AS attributed_orders,
    ROUND(SUM(revenue), 2)   AS attributed_revenue,     -- 精度: 小数2桁
    ROUND(SUM(revenue) / NULLIF(COUNT(DISTINCT order_id), 0), 2) AS avg_order_value
  FROM last_touch
  GROUP BY 1, 2
)

SELECT * FROM campaign_daily

stg_email_clicks.sql — ステージングモデル

-- models/staging/stg_email_clicks.sql
-- 修正⑤: テーブル名ハードコードを廃止し source() を使用
-- 修正①: 必要カラムのみ SELECT
SELECT
  campaign_id,
  user_id,
  DATE(event_timestamp, 'Asia/Tokyo') AS event_date,
  event_timestamp,
  email_subject,
  link_url
FROM {{ source('events', 'email_clicks') }}
WHERE
  -- パーティション剪定: 必ず DATE フィルタを付ける(BQ コスト削減)
  event_date >= '2024-01-01'

sources.yml — freshness 設定(修正⑥)

# models/sources.yml
version: 2

sources:
  - name: events
    database: my-gcp-project
    schema: events
    tables:
      - name: email_clicks
        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 (アラート発火)
        columns:
          - name: campaign_id
            tests: [not_null]
          - name: user_id
            tests: [not_null]
          - name: event_date
            tests: [not_null]

  - name: orders
    database: my-gcp-project
    schema: orders
    tables:
      - name: order_items
        description: "注文明細(Cloud Spanner → BQ レプリケーション)"
        loaded_at_field: _synced_at
        freshness:
          warn_after:  {count: 1, period: hour}
          error_after: {count: 3, period: hour}
        columns:
          - name: order_id
            tests: [not_null, unique]
          - name: user_id
            tests: [not_null]
          - name: revenue
            tests:
              - not_null
              - dbt_utils.accepted_range:
                  min_value: 0

DataDog freshness アラート統合(Argo Workflows)

# freshness_alert.py — dbt source freshness を DataDog メトリクスに送信
"""dbt source freshness を DataDog statsd メトリクスに変換して送信。

CronWorkflow: 毎時 dbt source freshness --output json → このスクリプトに渡す
"""
from __future__ import annotations
import json, os, sys
from datadog import initialize, statsd   # pip install datadog

initialize(statsd_host="localhost", statsd_port=8125)

def report_freshness(json_output: str) -> None:
    """dbt source freshness の JSON 出力を DataDog ゲージメトリクスに変換。

    Args:
        json_output: `dbt source freshness --output json` の標準出力

    Returns:
        None(DataDog に statsd.gauge を送信)
    """
    results = json.loads(json_output)
    for result in results.get("results", []):
        source_name: str = result["unique_id"]   # e.g. "source.my_project.events.email_clicks"
        status: str      = result["status"]      # "pass" | "warn" | "error"
        elapsed_hours: float = result.get("max_loaded_at_time_ago_in_s", 0) / 3600

        # DataDog メトリクス: dbt.source.freshness.elapsed_hours
        statsd.gauge(
            "dbt.source.freshness.elapsed_hours",
            elapsed_hours,
            tags=[
                f"source:{source_name}",
                f"status:{status}",
                f"env:{os.environ.get('DD_ENV', 'production')}",
            ],
        )

        if status in ("warn", "error"):
            statsd.event(
                title=f"dbt source freshness {status.upper()}: {source_name}",
                message=f"経過時間: {elapsed_hours:.1f}h — パイプライン遅延を確認してください",
                alert_type=status,
                tags=[f"source:{source_name}"],
            )

if __name__ == "__main__":
    report_freshness(sys.stdin.read())

mart_campaign_attribution_pivot.sql — BigQuery PIVOT(修正⑦)

-- models/mart/mart_campaign_attribution_pivot.sql
-- 修正⑦: pandas pivot_table → BigQuery PIVOT 演算子(Pod メモリ消費ゼロ)
-- キャンペーン別・月別の売上貢献を横持ち変換(直近6ヶ月)
{{
  config(
    materialized = 'table',   -- PIVOT 結果は table で完全リビルド
  )
}}

SELECT *
FROM (
  SELECT
    campaign_id,
    FORMAT_DATE('%Y-%m', click_date) AS ym,    -- 例: '2026-01'
    attributed_revenue
  FROM {{ ref('mart_campaign_attribution') }}
  WHERE click_date >= DATE_SUB(CURRENT_DATE('Asia/Tokyo'), INTERVAL 6 MONTH)
)
PIVOT (
  -- 修正⑦: BigQuery PIVOT 演算子で BQ クラスター上処理 → OOM なし
  -- 注意: dbt では静的 PIVOT のみサポート(動的 PIVOT は EXECUTE IMMEDIATE を使う)
  SUM(attributed_revenue) AS revenue
  FOR ym IN (
    '2025-12',
    '2026-01',
    '2026-02',
    '2026-03',
    '2026-04',
    '2026-05',
    '2026-06'
  )
)

動的 PIVOT(Jinja マクロ版)

-- macros/generate_pivot_months.sql
-- dbt Jinja で直近6ヶ月の月リストを動的生成
{% macro generate_pivot_months(n=6) %}
  {% set months = [] %}
  {% for i in range(n, 0, -1) %}
    {% set _ = months.append(
      "'" ~ (modules.datetime.date.today().replace(day=1) - modules.datetime.timedelta(days=30*i)).strftime('%Y-%m') ~ "'"
    ) %}
  {% endfor %}
  {{ months | join(', ') }}
{% endmacro %}

-- 使用例(mart_campaign_attribution_pivot.sql で呼び出し):
-- FOR ym IN ( {{ generate_pivot_months(6) }} )

dbt_project.yml — 変数定義(修正③)

# dbt_project.yml(抜粋)
name: mops_dbt
version: '1.0.0'
config-version: 2

# 修正③: マジックナンバーを vars で一元管理
vars:
  attribution_window_days: 7   # クリック後 N 日以内の注文を貢献とみなす

models:
  mops_dbt:
    staging:
      +materialized: view
    mart:
      +materialized: incremental
      +partition_by:
        field: click_date
        data_type: date
      +cluster_by: [campaign_id]

問題点と改善方法の対応表

問題修正内容効果
① SELECT *必要6カラムのみ明示スキャン85%削減($150 → $22.5/月)
② N+1 Python ループBigQuery JOIN 1本処理時間: 数時間 → 数十秒
③ マジックナンバー 7dbt var('attribution_window_days', 7)変更が dbt_project.yml 1行で完結
④ attribution 重複ROW_NUMBER() + WHERE click_rank=1ラスト・タッチ attribution が正確に
⑤ テーブル名ハードコード{{ source('events', 'email_clicks') }}依存グラフが有効・移行時に1箇所変更
⑥ freshness 未設定loaded_at_field + warn/error_afterパイプライン遅延を自動検知
⑦ pandas PIVOT OOMBigQuery PIVOT 演算子Pod メモリ消費ゼロ・OOMKilled なし

ポイント解説

1 BigQuery ウィンドウ関数 vs Python ループ(N+1 問題)
ROW_NUMBER() OVER(PARTITION BY order_id ORDER BY click_time DESC) は BigQuery の分散処理エンジンで並列実行される。Python での 100万行 × 1クエリ/行 = 100万 API 呼び出しは数時間かかる処理が数十秒に短縮される。BigQuery の設計原則は「データを移動させるのではなく、クエリをデータの近くで実行する」こと。Argo Workflows の Pod は処理ロジックの制御フローのみを担当し、データ処理は BigQuery に委譲するのがベストプラクティス。
2 ラスト・タッチ attribution の実装パターン
ECサイトでよく使われるアトリビューションモデルは3種類: (1) ラスト・タッチ: ROW_NUMBER() DESC = 1(最後のタッチポイントに100%帰属)、(2) ファースト・タッチ: ROW_NUMBER() ASC = 1(最初のタッチポイント)、(3) 線形: 1.0 / COUNT(*) OVER(PARTITION BY order_id)(均等配分)。モデルを変更する際は dbt の var() でパラメータ化しておくと設定変更のみで対応できる。
3 dbt source freshness の仕組みと DataDog 統合
dbt source freshnessloaded_at_field の最大値を取得し、現在時刻との差分(elapsed seconds)を計算する。warn_after: {count: 2, period: hour} は差分が 7,200秒を超えると warn ステータスになる。CronWorkflow で毎時実行し、--output json の結果を DataDog statsd に送ることで、既存の DataDog アラートルールと統合できる。Streaming Insert の場合は BQ システム列 _PARTITIONTIME または _ingested_at(自前で設定)を使う。
4 BigQuery PIVOT 演算子の制約と回避策
BigQuery の PIVOTFOR column IN (...) の値リストが静的(コンパイル時決定)なため、動的なカラム数には対応していない。動的 PIVOT が必要な場合は: (1) EXECUTE IMMEDIATE + テンプレート文字列で動的 SQL を生成、(2) dbt Jinja マクロで月リストを自動生成(コンパイル時)、の2パターンがある。dbt では Jinja マクロアプローチが実用的(BQ API 呼び出しなしで SQL 生成)。
5 BigQuery のスキャンコスト最適化: SELECT * vs 列指定
BigQuery はオンデマンドプランで $5/TB スキャン。列指向ストレージのため取得カラム数でコストが変わる。50カラムのうち6カラムのみ取得すると論理スキャン量が 6/50 = 12% になり、理論上 88% コスト削減。実際はカラムサイズのばらつきがあるため 85% 程度の削減が見込まれる。パーティション剪定(WHERE event_date >= ...)と組み合わせると更に大幅削減が可能。INFORMATION_SCHEMA.PARTITIONS で各パーティションの行数・サイズを事前確認すること。
6 dbt Mesh cross-project ref(dbt Core 1.8+)
複数チームが別 dbt プロジェクトを管理する場合、{{ ref('orders_project', 'stg_order_items') }} で別プロジェクトのモデルを参照できる(dbt Core 1.8+ / dbt Mesh)。orders.order_items を直接 source() で参照する代わりに Orders チームが公開した public model を ref() で使うと、スキーマ変更時に dbt が依存関係を自動で解決し、データ契約(Data Contract)が明示化される。

実務への応用

  • MOps attribution 分析の ROAS 計算: mart_campaign_attributionattributed_revenue / campaign_cost で ROAS(Return on Ad Spend)を算出。マーケターが「セグメント配信 vs 全体配信の ROAS 比較」を求めた場合は SUM(revenue) OVER(PARTITION BY segment_type) ウィンドウ関数で拡張できる。Looker Studio でキャンペーン別 ROAS ダッシュボードに直結させる。
  • dbt freshness を Slack アラートに連携: Argo Workflows の CronWorkflow で dbt source freshness --output json | python freshness_alert.py を毎時実行。DataDog モニターで dbt.source.freshness.elapsed_hours > 2 をトリガーに Slack の #mops-alerts チャンネルへ通知。「今朝の KPI が更新されていない」という手動確認業務をゼロにできる。
  • BigQuery INFORMATION_SCHEMA でパーティション健全性を確認: SELECT partition_id, total_rows, total_logical_bytes FROM \`my-project.events.INFORMATION_SCHEMA.PARTITIONS\` WHERE table_name = 'email_clicks' ORDER BY partition_id DESC LIMIT 7 で直近7日分のパーティションサイズを確認。突然の行数減少はパイプライン障害の早期シグナルになる。dbt テストとして dbt_utils.recency や custom generic test として組み込むことも可能。
  • attribution モデル切り替えのパラメータ化: dbt var で attribution_model: last_touch を管理し、Jinja で {% if var('attribution_model') == 'linear' %} の分岐を入れると、ファースト・タッチ / ラスト・タッチ / 線形の3モデルを同一 SQL で切り替えられる。マーケターからモデル変更を求められた際に SQL 改修不要でリリースできる。

今日のまとめ

Python ループで行ごとに BigQuery に問い合わせる N+1 パターンは「BigQuery JOIN + ROW_NUMBER() ウィンドウ関数」に置き換えるだけで処理時間を数時間から数十秒に短縮できる。

データエンジニアリングの設計原則は「データを Python に運ぶな、クエリをデータの近くで動かせ」。dbt では source() + freshness + var() の3点セットで依存管理・鮮度監視・マジックナンバー排除を同時に実現し、pandas PIVOT はメモリ OOM リスクがあるため BigQuery PIVOT 演算子に移行することが現代的なデータエンジニアリングの標準。

自己評価

自分の回答

気づき・メモ