【テクニカル・上級編】NotionとGoogle BigQueryのデータ連携:データベースの巨大ログをBIツールやSQLで高度に分析・可視化するパイプライン構築術 – プロジェクト・ナレッジ管理活用バイブル

NotionとGoogle BigQueryの完全同期パイプライン:巨大ログとナレッジをBIの海へ解き放つ極限アーキテクチャ

組織のスケールに伴い、Notionは単なるドキュメントツールから、プロジェクトのバックログ、意思決定のログ、リソース配分のマトリクスが混在する「ミッションクリティカルなデータベース」へと変貌を遂げる。しかし、ここでエンジニアリング組織が直面する壁は常に同じだ。

「Notion単体では、複雑な時系列集計、クロスデータベースの結合、そして高度なBI可視化が不可能である」という事実だ。

数万件を超えるタスクログ、膨大な工数データ、散在するナレッジのメタデータを人間の手や脆弱なZapierのワークフローでどうにかしようなどというアプローチは、今すぐ捨て去るべきだ。データはレプリカされ、構造化され、DWH(Data Warehouse)の土俵に上がって初めてビジネスの武器となる。

本稿では、NotionのREST APIの癖を骨の髄まで理解し、Google BigQuery(以下、BQ)へ完全自動かつ堅牢な差分同期パイプラインを構築するための、妥協なき設計思想と実装コードを提示する。

—

1. アーキテクチャ全体像:なぜ「直結」ではなく「パイプライン」が必要か

Notion APIは強力だが、リクエスト制限(Rate Limit:平均3リクエスト/秒)や、ネストされたプロパティ、Rich Textの厄介なJSON構造など、そのままBIツール(Looker StudioやTableauなど)のクエリソースとして耐えうる代物ではない。

我々が目指すべきは、以下の三層構造からなるパイプラインだ。

[Notion Database]
│
▼ (Incremental API Polling / Webhook + Queue)
[Extraction Layer (Python / Serverless or ECS)]
│
▼ (Schema Normalization & JSON Flattening)
[Google BigQuery (Raw & Staging Datasets)]
│
▼ (dbt / SQL Views)
[BI & Analytics (Looker Studio)]

このアーキテクチャの核心は、「Notionのスキーマ変更に強い正規化レイヤー」と「APIのレートリミットを完全制御するキューイング・リトライ戦略」の二点にある。

—

2. データの罠:Notion APIの深淵と向き合う

Notionからデータを引き抜く際、エンジニアが直面する最大の障害は「スキーマの流動性」と「データの非正規化」である。

プロパティのフラット化問題

Notionのデータベースのレスポンスは、以下のようにプロパティごとに型がラップされた冗長なJSONオブジェクトを返す。

“Status”: {
“id”: “abc1”,
“type”: “status”,
“status”: {
“id”: “xyz2”,
“name”: “In Progress”,
“color”: “blue”
}
}

これをBQへそのまま突っ込んでも、SQLでの集計(`GROUP BY` や `JOIN`)で地獄を見る。抽出レイヤーにおいて、このJSONを以下のようなフラットなスキーマへ変換(Transformation)することが絶対条件となる。

  • `page_id` (STRING)
  • `created_time` (TIMESTAMP)
  • `last_edited_time` (TIMESTAMP)
  • `title` (STRING)
  • `status` (STRING)
  • `assignee_ids` (ARRAY)

—

3. 実装:Pythonによる堅牢な差分同期スクリプト

ここでは、サードパーティの重いSaaS(AirbyteやFivetran)を使わず、自社でコントロール可能な極限まで最適化されたPython製抽出・ロードスクリプトの核心部分を解説する。

Notion APIの `last_edited_time` を活用し、前回同期以降に更新されたページのみを取得する差分同期(Incremental Sync)を実装する。

差分同期スクリプト (`notion_to_bq_sync.py`)

import os
import time
from datetime import datetime, timezone
from google.cloud import bigquery
from notion_client import Client, APIResponseError

環境変数からの設定読み込み
NOTION_TOKEN = os.environ[“NOTION_TOKEN”]
NOTION_DATABASE_ID = os.environ[“NOTION_DATABASE_ID”]
BQ_PROJECT_ID = os.environ[“BQ_PROJECT_ID”]
BQ_DATASET_ID = os.environ[“BQ_DATASET_ID”]
BQ_TABLE_ID = os.environ[“BQ_TABLE_ID”]

notion = Client(auth=NOTION_TOKEN)
bq_client = bigquery.Client(project=BQ_PROJECT_ID)

def get_last_sync_timestamp() -> str:
“””BigQuery上の対象テーブルから最新のlast_edited_timeを取得し、
APIコールの下限とする(なければ十分古い日時を返す)。”””
table_ref = f”{BQ_PROJECT_ID}.{BQ_DATASET_ID}.{BQ_TABLE_ID}”
try:
query = f”SELECT MAX(last_edited_time) as max_time FROM `{table_ref}`”
query_job = bq_client.query(query)
result = list(query_job.result())
if result and result[0][“max_time”]:
# ISO8601形式へ変換
return result[0][“max_time”].astimezone(timezone.utc).isoformat()
except Exception as e:
print(f”[Warning] テーブルが存在しないか、初期ロードとみなします: {e}”)

# デフォルト:1年前(初回フル同期に近い挙動)
return “2023-01-01T00:00:00.000Z”

def fetch_notion_pages(start_cursor=None, filter_timestamp=None):
“””Notion APIを叩いてページをページネーション付きで取得。
Rate Limit (3 req/sec) を考慮し、ウェイトを入れる。”””
query_params = {
“database_id”: NOTION_DATABASE_ID,
“page_size”: 100,
}
if start_cursor:
query_params[“start_cursor”] = start_cursor

if filter_timestamp:
query_params[“filter”] = {
“timestamp”: “last_edited_time”,
“last_edited_time”: {
“after”: filter_timestamp
}
}
query_params[“sorts”] = [{
“timestamp”: “last_edited_time”,
“direction”: “ascending”
}]

try:
time.sleep(0.34) # 3 req/secの制限を遵守するためのスロットリング
return notion.databases.query(query_params)
except APIResponseError as e:
if e.status == 429:
print(“[Error] Rate limit exceeded. 60秒待機します…”)
time.sleep(60)
return fetch_notion_pages(start_cursor, filter_timestamp)
raise e

def transform_page(page: dict) -> dict:
“””Notionの冗長なJSONオブジェクトを、BQに最適化されたフラットな構造へパースする”””
props = page.get(“properties”, {})

# 例:Titleプロパティの抽出
title_prop = props.get(“Name”, {}).get(“title”, [])
title_text = “”.join([t.get(“plain_text”, “”) for t in title_prop])

# 例:Statusプロパティの抽出
status_prop = props.get(“Status”, {}).get(“status”)
status_name = status_prop.get(“name”) if status_prop else None

# 例:People(Assignee)プロパティの抽出
assignees = props.get(“Assignee”, {}).get(“people”, [])
assignee_ids = [a.get(“id”) for a in assignees]

return {
“page_id”: page.get(“id”),
“created_time”: page.get(“created_time”),
“last_edited_time”: page.get(“last_edited_time”),
“archived”: page.get(“archived”),
“title”: title_text,
“status”: status_name,
“assignee_ids”: assignee_ids,
“raw_json”: str(page) # デバッグおよび将来のスキーマ拡張用の生JSON保存
}

def load_to_bigquery(rows: list):
“””BQのストリーミングインサート、またはロードジョブでデータを投入”””
if not rows:
print(“投入するデータがありません。”)
return

table_id = f”{BQ_PROJECT_ID}.{BQ_DATASET_ID}.{BQ_TABLE_ID}”

# 差分更新(Upsert)を実現するため、一時テーブルへロードしてMERGEする戦略をとるのがベスト
temp_table_id = f”{table_id}_temp_{int(time.time())}”

job_config = bigquery.LoadJobConfig(
autodetect=True,
write_disposition=bigquery.WriteDisposition.WRITE_TRUNCATE,
)

# 一時テーブルに書き込み
load_job = bq_client.load_table_from_json(rows, temp_table_id, job_config=job_config)
load_job.result() # 完了を待つ

# MERGEクエリで本番テーブルへUpsert
merge_query = f”””
MERGE `{table_id}` T
USING `{temp_table_id}` S
ON T.page_id = S.page_id
WHEN MATCHED THEN
UPDATE SET
T.last_edited_time = S.last_edited_time,
T.archived = S.archived,
T.title = S.title,
T.status = S.status,
T.assignee_ids = S.assignee_ids,
T.raw_json = S.raw_json
WHEN NOT MATCHED THEN
INSERT (page_id, created_time, last_edited_time, archived, title, status, assignee_ids, raw_json)
VALUES (page_id, created_time, last_edited_time, archived, title, status, assignee_ids, raw_json);
“””
bq_client.query(merge_query).result()

# 一時テーブルの削除
bq_client.delete_table(temp_table_id)
print(f”Successfully upserted {len(rows)} rows into BigQuery.”)

def main():
last_sync = get_last_sync_timestamp()
print(f”Syncing changes since: {last_sync}”)

has_more = True
start_cursor = None
all_transformed_rows = []

while has_more:
response = fetch_notion_pages(start_cursor=start_cursor, filter_timestamp=last_sync)
results = response.get(“results”, [])

for page in results:
all_transformed_rows.append(transform_page(page))

has_more = response.get(“has_more”, False)
start_cursor = response.get(“next_cursor”)

load_to_bigquery(all_transformed_rows)

if __name__ == “__main__”:
main()

—

4. パフォーマンスとメモリ消費の最適化ハック

上記のようなスクリプトをKubernetes(CronJob)やCloud Run Jobsで運用する際、データ量が数万件〜数十万件規模にスケールすると、いくつかの深刻なボトルネックに突き当たる。

1. バッチ処理によるメモリ枯渇(OOM)の回避

すべてのページをメモリ上に配列として蓄積してからBQへ一括ロードすると、ページネーションが数千回に及んだ際にコンテナがメモリ不足(OOM Killer)で強制終了する。
対策: 蓄積バッチサイズ(例:500件ごと)に達した時点で一度BQへロード処理(Flush)を挟むストリーミング・バッチ方式へリファクタリングせよ。

2. BigQueryのパーティショニングとクラスタリング

BQに作成するテーブルは、無思考に作るのではなく、以下のDDLを適用してストレージコストとスキャンコストを極限まで削れ。

CREATE TABLE `your_project.your_dataset.notion_tasks`
(
page_id STRING NOT NULL,
created_time TIMESTAMP,
last_edited_time TIMESTAMP,
archived BOOLEAN,
title STRING,
status STRING,
assignee_ids ARRAY,
raw_json STRING
)
PARTITION BY DATE(last_edited_time)
CLUSTER BY status, page_id;

  • パーティショニング (`PARTITION BY`): タイムスタンプベースでパーティションを分けることで、BIツールから直近のデータのみをスキャンさせ、クエリコストを劇的に削減する。
  • クラスタリング (`CLUSTER BY`): 頻繁にフィルタリングやグルーピングの条件になる `status` などをクラスタ化キーに指定し、I/Oを最適化する。

—

5. 経営・開発ダッシュボードの構築とネクストステップ

BigQueryへのデータパイプラインが完成すれば、あとはLooker StudioやTableauの出番だ。

  • ベロシティの時系列推移: ステータスが「Done」になったタスクの数をタイムスタンプで日別・週別に集計し、スループットを可視化。
  • ボトルネック分析: 各タスクがどのステータス(Review, In Progress等)に何日滞留しているかを `last_edited_time` の差分から算出し、開発プロセスの詰まりを高精度に検知する。

データエンジニアとしての美しさは、「人が手作業で行う情報の同期や集計を、コードによって完全に無力化する」ことにある。Notionという情報の源泉をBigQueryという大河に繋ぎ、組織全体の意思決定スピードを物理的限界まで引き上げてほしい。

タイトルとURLをコピーしました