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という大河に繋ぎ、組織全体の意思決定スピードを物理的限界まで引き上げてほしい。