概要
相関サブクエリを window 関数で1パス化
「前のイベント時刻」を SELECT MAX(...) WHERE p.time < e.time の相関サブクエリで取ると各行ごとに再スキャンが走り実質 O(N²) になる。LAG(event_time) OVER (PARTITION BY user ORDER BY time) は BigQuery が1回のソートで解決し、スキャン量・スロット消費・実行時間が桁違いに小さくなる。
gaps-and-islands でセッション採番
セッション化の定石。新セッション開始フラグ(前イベントとの間隔が30分超、または最初の行で LAG が NULL)を 1/0 で作り、その 累積和 SUM(flag) OVER (... ROWS UNBOUNDED PRECEDING) をセッション連番として使う。user_id || '-' || seq で一意な session_id を得る。
ファネルは条件付き集約で1パス
view→cart→checkout→purchase の到達判定をステップごとに4回 self-join すると重い。セッション単位に GROUP BY し、LOGICAL_OR(event_type='add_to_cart') や COUNTIF、順序を見るなら MIN(IF(...)) で1パス集約。テーブルスキャンは1回で済む。
近似カウント + 増分 + contract
巨大 COUNT(DISTINCT) は APPROX_COUNT_DISTINCT(HLL)でメモリ/コスト削減。materialized: table は incremental + insert_overwrite + partition/cluster でパーティション洗い替え。QUALIFY で ROW_NUMBER サブクエリを排除。dbt model contract でスキーマドリフトを build 時に fail。
問題
ECサイト MOps チームでは、Webフロントのユーザー行動ログ(analytics.pageview_events、1行=1イベント)を BigQuery × dbt Core 1.9+ で「セッション別ファネル集計マート」に変換している。目的は、キャンペーン流入ユーザーの view → add_to_cart → checkout → purchase ファネル通過率を、30分アイドルで区切ったセッション単位 で日次計測すること。以下の dbt モデル・スキーマには 7つの設計・パフォーマンス・データ品質上の問題 が潜んでいる。問題点を全て洗い出し、window 関数(LAG / gaps-and-islands / QUALIFY)× 条件付き集約 × APPROX_COUNT_DISTINCT × incremental insert_overwrite + partition/cluster × dbt model contract を活用した Bad→Good リファクタリングを行え。
制約・前提条件
- BigQuery Standard SQL(GA 機能のみ。
QUALIFY,APPROX_COUNT_DISTINCT,HLL_COUNT.*,COUNTIF,LOGICAL_OR利用可) - dbt Core 1.9+(model contract
enforced: true、constraints、incremental_strategy='insert_overwrite'利用可) - テーブル
analytics.pageview_events(partition:event_date DATE, cluster:user_id)。日次で数億行、直近1日分のみ増分 - セッション定義: 同一
user_id内でイベント間隔が 30分以上空いたら新セッション - BigQuery オンデマンド課金($5/TB)。相関サブクエリ・自己結合の多用でスキャン量・スロット消費が過大
- 下流マートがこのモデルの
session_idを join キーに使うため、NULL や重複が混入するとサイレントに壊れる
悪いコード (Before)
-- 問題⑥: table = 毎回フルスキャン、partition/cluster なし
{{ config(materialized='table') }}
WITH events AS (
SELECT user_id, event_time, event_type, channel, event_date
FROM {{ source('analytics','pageview_events') }}
),
with_prev AS (
SELECT
e.user_id, e.event_time, e.event_type, e.channel,
-- 問題①: 相関サブクエリで前イベント時刻取得 → O(N^2)
(SELECT MAX(p.event_time)
FROM events p
WHERE p.user_id = e.user_id
AND p.event_time < e.event_time) AS prev_time
FROM events e
),
-- 問題②: セッション境界を採番せず、日付でざっくり集計
-- → 1日に複数回来訪しても1セッション扱い
ranked AS (
SELECT *,
ROW_NUMBER() OVER (PARTITION BY user_id, event_date
ORDER BY event_time) AS rn
FROM with_prev
),
-- 問題③: ROW_NUMBER をサブクエリで包んで WHERE rn=1
first_ev AS ( SELECT * FROM ranked WHERE rn = 1 ),
-- 問題④: ファネル4ステップを4回自己結合
steps AS ( SELECT user_id, event_date, event_type FROM with_prev ),
funnel AS (
SELECT v.user_id, v.event_date
FROM steps v
LEFT JOIN steps c ON c.user_id=v.user_id AND c.event_date=v.event_date
AND c.event_type='add_to_cart'
LEFT JOIN steps k ON k.user_id=v.user_id AND k.event_date=v.event_date
AND k.event_type='checkout'
LEFT JOIN steps p ON p.user_id=v.user_id AND p.event_date=v.event_date
AND p.event_type='purchase'
WHERE v.event_type='view'
)
SELECT
event_date,
-- 問題⑤: 数億行に対する COUNT(DISTINCT) が高コスト
COUNT(DISTINCT user_id) AS users
FROM funnel
GROUP BY event_date
version: 2
models:
- name: bad_session_funnel
# 問題⑦: contract なし・constraints なし
# → 上流のカラム追加/型変更で session_id が
# NULL/重複してもサイレントに下流マートを破壊
columns:
- name: session_id # そもそも生成していない
- name: event_date
LAG(event_time) OVER (PARTITION BY user ORDER BY time) で1パス化session_id 採番WHERE rn=1 の入れ子。QUALIFY ROW_NUMBER()=1 で1文にLOGICAL_OR/COUNTIF の条件付き集約で1パスAPPROX_COUNT_DISTINCT(HLL)で近似incremental + insert_overwrite + partition/clustermodel contract enforced + constraintsヒント(段階的開示)
ヒント1 — 方向性
SELECT MAX(...) WHERE p.event_time < e.event_time)で取ると実質 O(N²) スキャンになる。これは LAG(event_time) OVER (PARTITION BY user_id ORDER BY event_time) の1パスで置き換えられる。セッションIDの採番は gaps-and-islands パターン、すなわち「新セッション開始フラグ(間隔30分超)」を作り、その 累積和(SUM(flag) OVER (...))をセッション連番として使う定石で解ける。「各セッションの最初のイベントだけ」を取るのに ROW_NUMBER をサブクエリで包んで WHERE rn=1 する必要はなく、QUALIFY で1文にできる。
ヒント2 — アプローチ
- 問題①: 相関サブクエリ(O(N²))→
LAG(event_time) OVER (PARTITION BY user_id ORDER BY event_time) - 問題②: セッション未採番 → gaps-and-islands(
is_new_sessionの累積SUMでsession_id) - 問題③:
ROW_NUMBERサブクエリ+WHERE rn=1→QUALIFY - 問題④: ファネル4回自己結合 → セッション単位の条件付き集約(
COUNTIF/LOGICAL_OR/MIN(IF())) - 問題⑤: 巨大
COUNT(DISTINCT)→APPROX_COUNT_DISTINCT(必要ならHLL_COUNT.INIT/MERGE) - 問題⑥:
table→incremental+insert_overwrite+partition_by(event_date)+cluster_by+ パーティション剪定 - 問題⑦: スキーマ無防備 → dbt
model contract enforced: true+constraints(not_null / primary_key)
ヒント3 — コードの骨格
-- gaps-and-islands セッション化の骨格
WITH flagged AS (
SELECT
user_id, event_time, event_type, event_date,
TIMESTAMP_DIFF(
event_time,
LAG(event_time) OVER (PARTITION BY user_id ORDER BY event_time),
MINUTE
) AS gap_min -- 修正①: 前イベント間隔
FROM {{ ref('stg_pageview_events') }}
),
sessionized AS (
SELECT *,
-- 修正②: 30分超 or 最初(gap NULL) → 新セッション=1、その累積和が連番
SUM(CASE WHEN gap_min IS NULL OR gap_min >= 30 THEN 1 ELSE 0 END)
OVER (PARTITION BY user_id ORDER BY event_time
ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW) AS session_seq
FROM flagged
)
-- session_id = user_id || '-' || session_seq
-- 以降はセッション単位で LOGICAL_OR / COUNTIF 集約(修正④)
-- 流入チャネルは QUALIFY ROW_NUMBER()=1 で取得(修正③)
問題点分析(7点)
| # | 問題点 | 分類 | 改善方法 |
|---|---|---|---|
| 1 | 相関サブクエリで前イベント時刻取得(O(N²)) | パフォーマンス | LAG() window 関数で1パス化 |
| 2 | セッション未採番(1日複数来訪を1セッション扱い) | 正確性 | gaps-and-islands(新セッションフラグの累積 SUM) |
| 3 | ROW_NUMBER をサブクエリ+WHERE rn=1 | 可読性/効率 | QUALIFY で1文に |
| 4 | ファネル4ステップを4回自己結合 | パフォーマンス | COUNTIF / LOGICAL_OR の条件付き集約1パス |
| 5 | 巨大 COUNT(DISTINCT user_id) | コスト/メモリ | APPROX_COUNT_DISTINCT(HLL) |
| 6 | materialized: table(フルスキャン・無 partition/cluster) | コスト | incremental insert_overwrite + partition/cluster |
| 7 | contract なし(型/NULL/重複を検知不能) | データ品質 | model contract enforced + constraints |
セッション化の仕組み — gaps-and-islands(LAG + 累積 SUM)
模範解答
{{ config(materialized='table') }}
-- 前イベント時刻を相関サブクエリで(O(N^2))
(SELECT MAX(p.event_time) FROM events p
WHERE p.user_id=e.user_id
AND p.event_time < e.event_time) AS prev_time
-- セッション未採番、日付でざっくり集計
-- ROW_NUMBER をサブクエリで包み WHERE rn=1
-- ファネルは4回 self-join
-- COUNT(DISTINCT user_id) が高コスト
-- models/marts/mart_session_funnel.sql | dbt Core 1.9+
{{ config(
materialized='incremental',
incremental_strategy='insert_overwrite', -- 修正⑥: 冪等な洗い替え
partition_by={'field':'event_date',
'data_type':'date','granularity':'day'},
cluster_by=['first_channel'],
on_schema_change='fail'
) }}
WITH source_events AS (
SELECT user_id, session_user_key, event_time,
event_type, channel, event_date
FROM {{ ref('stg_pageview_events') }}
WHERE event_type IN ('view','add_to_cart','checkout','purchase')
{% if is_incremental() %}
-- 修正⑥: パーティション剪定(対象日のみ読む)
AND event_date = DATE('{{ var("run_date",
run_started_at.strftime("%Y-%m-%d")) }}')
{% endif %}
),
-- 修正①: 相関サブクエリ → LAG で前イベント間隔を1パス
flagged AS (
SELECT *,
TIMESTAMP_DIFF(event_time,
LAG(event_time) OVER (PARTITION BY session_user_key
ORDER BY event_time),
MINUTE) AS gap_min
FROM source_events
),
-- 修正②: gaps-and-islands でセッション連番を採番
sessionized AS (
SELECT *,
SUM(CASE WHEN gap_min IS NULL OR gap_min >= 30
THEN 1 ELSE 0 END)
OVER (PARTITION BY session_user_key ORDER BY event_time
ROWS BETWEEN UNBOUNDED PRECEDING
AND CURRENT ROW) AS session_seq
FROM flagged
),
with_session_id AS (
SELECT *,
CONCAT(session_user_key,'-',
CAST(session_seq AS STRING)) AS session_id -- 下流 join キー
FROM sessionized
),
-- 修正④: 4回自己結合 → セッション単位の条件付き集約(1パス)
session_funnel AS (
SELECT
session_id, session_user_key,
MIN(event_date) AS event_date,
MIN(event_time) AS session_start_at,
LOGICAL_OR(event_type='view') AS reached_view,
LOGICAL_OR(event_type='add_to_cart')AS reached_cart,
LOGICAL_OR(event_type='checkout') AS reached_checkout,
LOGICAL_OR(event_type='purchase') AS reached_purchase,
COUNTIF(event_type='purchase') AS purchase_events,
COUNT(*) AS event_count
FROM with_session_id
GROUP BY session_id, session_user_key
)
-- 修正③: 流入チャネルは QUALIFY で「セッション最初の1行」を取得
SELECT
f.event_date, c.first_channel, f.session_id,
f.session_user_key, f.session_start_at,
f.reached_view, f.reached_cart,
f.reached_checkout, f.reached_purchase,
f.purchase_events, f.event_count
FROM session_funnel f
JOIN (
SELECT session_id, channel AS first_channel
FROM with_session_id
QUALIFY ROW_NUMBER() OVER (PARTITION BY session_id
ORDER BY event_time) = 1 -- 修正③
) c USING (session_id)
-- models/marts/mart_daily_funnel_rate.sql
-- セッション mart を集約し、チャネル別の日次ファネル通過率を出す
SELECT
event_date,
first_channel,
-- 修正⑤: 巨大 DISTINCT → 近似カウント(誤差 ~1.6%、メモリ/コスト大幅減)
APPROX_COUNT_DISTINCT(session_user_key) AS approx_users,
COUNT(DISTINCT session_id) AS sessions, -- session は元々少数
COUNTIF(reached_view) AS f_view,
COUNTIF(reached_cart) AS f_cart,
COUNTIF(reached_checkout) AS f_checkout,
COUNTIF(reached_purchase) AS f_purchase,
SAFE_DIVIDE(COUNTIF(reached_cart), COUNTIF(reached_view)) AS cvr_view_to_cart,
SAFE_DIVIDE(COUNTIF(reached_purchase), COUNTIF(reached_view)) AS cvr_view_to_purchase
FROM {{ ref('mart_session_funnel') }}
-- パーティション剪定を必ず効かせる
WHERE event_date = DATE('{{ var("run_date", run_started_at.strftime("%Y-%m-%d")) }}')
GROUP BY event_date, first_channel
HLL_COUNT.INIT(session_user_key) のスケッチ列を保存し、期間横断で HLL_COUNT.MERGE すると重複排除カウントを再スキャンなしで再集計できる。
# models/marts/_marts.yml
# 修正⑦: contract で型・NULL・一意性を「モデル build 時」に強制
version: 2
models:
- name: mart_session_funnel
description: "セッション別ファネル到達フラグ(30分アイドルで区切ったセッション単位)"
config:
contract:
enforced: true # スキーマ不一致・制約違反で build を fail
columns:
- name: session_id
data_type: string
constraints:
- type: not_null # 下流 join キー: NULL 混入を build 時に検知
- type: primary_key # 一意性の意図を明示(BQ は informational)
tests:
- unique # 実一意性は dbt test で検証(gaps-and-islands の正しさ担保)
- name: event_date
data_type: date
constraints:
- type: not_null
- name: first_channel
data_type: string
tests:
- accepted_values:
values: ['organic','paid_search','email','social','direct']
- name: reached_view
data_type: boolean
constraints:
- type: not_null
- name: reached_purchase
data_type: boolean
constraints:
- type: not_null
contract.enforced: true は宣言した data_type と実出力の不一致で build を落とす。Web計測タグ改修でイベントスキーマが変わっても、session_id の型崩れ・NULL 化を下流に流す前に止められる。
| 問題 | 修正内容 | 効果 |
|---|---|---|
| ① 相関サブクエリ | LAG() OVER window | O(N²) → 1パス。スキャン・スロット・実行時間を大幅削減 |
| ② セッション未採番 | gaps-and-islands(新セッションフラグの累積 SUM) | 1日複数来訪を正しく分割、通過率を過小評価しない |
| ③ ROW_NUMBER サブクエリ | QUALIFY ROW_NUMBER()=1 | 入れ子を排除、1文で「最初の行」抽出 |
| ④ 4回自己結合 | LOGICAL_OR / COUNTIF 条件付き集約 | テーブルスキャン1回、結合を消去 |
| ⑤ COUNT(DISTINCT) | APPROX_COUNT_DISTINCT(HLL) | 誤差 ~1-2% でメモリ/コスト大幅減、スケッチ再利用可 |
| ⑥ table フルスキャン | incremental insert_overwrite + partition/cluster | 対象日パーティションのみ洗い替え(冪等)+ 剪定 |
| ⑦ contract なし | model contract enforced + constraints | スキーマドリフト/NULL/型崩れを build 時に fail fast |
実行例(input → output)
user_id event_time event_type channel
U1 2026-07-15 10:00:00 view email
U1 2026-07-15 10:05:00 add_to_cart email
U1 2026-07-15 10:20:00 checkout email ← 15分: 同一セッション
U1 2026-07-15 11:30:00 view email ← 70分: 新セッション
U1 2026-07-15 11:32:00 purchase email
session_id first_channel view cart checkout purchase
U1-1 email true true true false
U1-2 email true false false true
LAG の gap が 70分 > 30分 → session_seq が 1→2 にインクリメントされ正しく2分割ポイント解説
SELECT MAX(...) WHERE p.time < e.time で取ると各行ごとに再スキャンが走り実質 O(N²)。LAG(col) OVER (PARTITION BY key ORDER BY time) は BigQuery が1回のソートで解決するため、スキャン量・スロット消費・実行時間が桁違いに小さい。「次の行」は LEAD。LAG が NULL)を 1/0 で作り、(b) その累積和 SUM(flag) OVER (PARTITION BY user ORDER BY time ROWS UNBOUNDED PRECEDING) を取ると、フラグが立つたびに連番が +1 され、これがそのままセッション連番になる。user_id || '-' || seq で一意な session_id を得る。SELECT * FROM (SELECT *, ROW_NUMBER() OVER(...) rn) WHERE rn=1 と入れ子にする代わりに、QUALIFY ROW_NUMBER() OVER(...) = 1 と1文で書ける。window 関数の結果に対する WHERE に相当し、可読性が上がりオプティマイザにも優しい。GROUP BY し、各ステップは LOGICAL_OR(event_type='add_to_cart')(1回でも起きたか)や COUNTIF(...)、順序を見たいなら MIN(IF(event_type='checkout', event_time, NULL)) で到達時刻を1パスで得る。テーブルスキャンは1回で済む。COUNT(DISTINCT) は全カーディナリティをメモリ保持するため数億ユーザーで高コスト。APPROX_COUNT_DISTINCT(HyperLogLog++)は誤差 ~1-2% と引き換えにメモリ・実行コストを大幅削減。日次スケッチを再利用したいなら HLL_COUNT.INIT でスケッチ列を作り、HLL_COUNT.MERGE で期間横断の重複排除カウントを安価に得られる(週次/月次 UU を日次スケッチから再集計)。materialized: table(毎回フルスキャン)を incremental に変え、insert_overwrite 戦略で「対象パーティション(event_date)だけを丸ごと洗い替え」する。再実行しても対象日パーティションが置換されるだけなので冪等(microbatch と並ぶ増分戦略で、パーティション単位で明示的に置換したい時に使う)。cluster_by で頻出フィルタ列のブロック剪定、クエリ側 WHERE event_date = ... でパーティション剪定を必ず効かせる。contract.enforced: true を付けると、columns に宣言した data_type と実出力スキーマが不一致だと build 時に失敗する。constraints(not_null / primary_key)は下流 join キーの品質を型レベルで守る(BigQuery で primary_key は informational だが、意図の明示 + unique テスト併用で実効性を持たせる)。上流のカラム追加・型変更をサイレントに通さず、session_id の NULL/重複といったマート破壊を build ゲートで止める。実務への応用
ECサイト MOps のキャンペーン流入 → 購入ファネル分析に直結する。セッション化を gaps-and-islands で正しく行うことで、「同じユーザーが1日に何度も来訪した」ケースを1セッションに潰さず、来訪ごとの通過率を計測できる(キャンペーンメール経由の再来訪を過小評価しない)。相関サブクエリ→LAG、4回自己結合→条件付き集約の置き換えだけで、数億行スキャンのクエリコスト・実行時間が大幅に下がり、オンデマンド課金を圧縮できる。
APPROX_COUNT_DISTINCT / HLL_COUNT スケッチは、日次で作ったスケッチを週次・月次 UU に MERGE で再利用でき、Looker Studio のダッシュボードで期間を切り替えても安価に近似 UU を出せる。
session_id の NULL/型崩れを build 時に検知するのが有効。
今日のまとめ
次のステップ
- 発展問題: このファネルを「ステップ順序を厳密化」して計測せよ(cart の後の checkout のみ有効、逆流は無効)。
MIN(IF(event_type='view', event_time, NULL))等で各ステップ最初到達時刻を取り、t_view <= t_cart <= t_checkout <= t_purchaseの単調増加を満たすセッションのみ「正規ファネル通過」と判定するロジックを設計せよ。さらに、離脱ステップ別の内訳(どのステップで何%が落ちたか)を1クエリで出すこと。 - 参考: BigQuery Window(Navigation: LAG/LEAD, Numbering: ROW_NUMBER)/ QUALIFY 句 / APPROX_COUNT_DISTINCT・HLL_COUNT.INIT/MERGE / gaps-and-islands パターン / dbt incremental strategies(insert_overwrite)/ dbt model contracts & constraints