概要
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 メトリクスへの鮮度アラートも設定すること
悪いコード (Before)
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 未設定 → データ遅延を検知できない
SELECT * によるフルカラムスキャン — 50カラム全取得で月額 $150。必要6カラムのみ取得で $22.5(85%削減)7 の散在 — SQL と Python の2箇所にハードコード。dbt var() で一元管理ヒント(段階的開示)
ヒント1 — 方向性
{{ 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点)
| # | 問題点 | 分類 | 改善方法 |
|---|---|---|---|
| 1 | SELECT * フルカラムスキャン | コスト | 必要6カラムのみ明示(85%削減) |
| 2 | Python N+1 BigQuery API ループ | パフォーマンス | BigQuery JOIN 1本に集約 |
| 3 | マジックナンバー 7 の散在 | 保守性 | dbt var('attribution_window_days', 7) |
| 4 | ラスト・タッチ attribution 未定義 | 正確性 | ROW_NUMBER() OVER + WHERE click_rank=1 |
| 5 | dbt source() 未使用・ハードコード | 依存管理 | {{ source('events', 'email_clicks') }} |
| 6 | dbt source freshness 未設定 | 可観測性 | loaded_at_field + warn/error_after |
| 7 | pandas PIVOT → Pod OOM | 信頼性 | BigQuery PIVOT 演算子 |
設計図 — Bad vs Good データフロー
模範解答
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本 | 処理時間: 数時間 → 数十秒 |
| ③ マジックナンバー 7 | dbt 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 OOM | BigQuery PIVOT 演算子 | Pod メモリ消費ゼロ・OOMKilled なし |
ポイント解説
ROW_NUMBER() OVER(PARTITION BY order_id ORDER BY click_time DESC) は BigQuery の分散処理エンジンで並列実行される。Python での 100万行 × 1クエリ/行 = 100万 API 呼び出しは数時間かかる処理が数十秒に短縮される。BigQuery の設計原則は「データを移動させるのではなく、クエリをデータの近くで実行する」こと。Argo Workflows の Pod は処理ロジックの制御フローのみを担当し、データ処理は BigQuery に委譲するのがベストプラクティス。
ECサイトでよく使われるアトリビューションモデルは3種類: (1) ラスト・タッチ:
ROW_NUMBER() DESC = 1(最後のタッチポイントに100%帰属)、(2) ファースト・タッチ: ROW_NUMBER() ASC = 1(最初のタッチポイント)、(3) 線形: 1.0 / COUNT(*) OVER(PARTITION BY order_id)(均等配分)。モデルを変更する際は dbt の var() でパラメータ化しておくと設定変更のみで対応できる。
dbt source freshness は loaded_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(自前で設定)を使う。
BigQuery の
PIVOT は FOR column IN (...) の値リストが静的(コンパイル時決定)なため、動的なカラム数には対応していない。動的 PIVOT が必要な場合は: (1) EXECUTE IMMEDIATE + テンプレート文字列で動的 SQL を生成、(2) dbt Jinja マクロで月リストを自動生成(コンパイル時)、の2パターンがある。dbt では Jinja マクロアプローチが実用的(BQ API 呼び出しなしで SQL 生成)。
BigQuery はオンデマンドプランで $5/TB スキャン。列指向ストレージのため取得カラム数でコストが変わる。50カラムのうち6カラムのみ取得すると論理スキャン量が 6/50 = 12% になり、理論上 88% コスト削減。実際はカラムサイズのばらつきがあるため 85% 程度の削減が見込まれる。パーティション剪定(
WHERE event_date >= ...)と組み合わせると更に大幅削減が可能。INFORMATION_SCHEMA.PARTITIONS で各パーティションの行数・サイズを事前確認すること。
複数チームが別 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_attributionのattributed_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 に運ぶな、クエリをデータの近くで動かせ」。dbt では
source() + freshness + var() の3点セットで依存管理・鮮度監視・マジックナンバー排除を同時に実現し、pandas PIVOT はメモリ OOM リスクがあるため BigQuery PIVOT 演算子に移行することが現代的なデータエンジニアリングの標準。