【実務・中級編】 Change Streamsの内部アーキテクチャ – Cloud Spanner

Spanner Change Streamsの内部アーキテクチャ:分散RDBの極限レイヤーをハックする

こんにちは。チーフアーキテクトの私だ。
今日のコードレビューや設計レビューで、こんな質問を受けていないだろうか?

  • 「Change Streamsって、要するにRDBのCDC(変更データキャプチャ)機能でしょ?Kafka Connectみたいな感覚で使って大丈夫?」
  • 「リアルタイムで全件ストリーミングしたいんだけど、パフォーマンス影響はないの?」
  • 「スプリットやマイグレーションが走ったとき、変更ストリームの順序保証や欠損はどう担保されているの?」

もし、これらに「なんとなく動きます」としか答えられないなら、君のアーキテクチャ設計は危険だ。
Cloud Spannerの Change Streams は、単なる「便利なトリガーの代替」ではない。グローバル分散データベースの深部、パーシステンス層とコンセンサス層の極限の最適化の上に成り立っている「分布式トランザクションログの直接購読インターフェース」なのだ。

今日は、その内部アーキテクチャの泥臭い仕組みから、実務で絶対に踏み抜いてはいけない設計パターンまで、一切の妥協なく解説しよう。

—

1. 内部アーキテクチャ:変更データはどこから生まれ、どう流れるのか

まず、Spannerのストレージとコンセンサスの基礎を思い出してほしい。
Spannerのデータは、キー範囲ごとに「スプリット(Split)」という単位に分割され、それぞれが複数のPaxonsグループ(レプリカ)によって管理されている。

Change Streamsが有効化されたテーブルで何が起きているか。内部のデータフロードライバは、以下のようなレイヤーで動いている。

[Client Mutation]
↓
[Paxos Leader (Split)] —> (Distributed Commit)
↓
[RocksDB (Storage Layer)]
↓ (背景で非同期スキャン)
[Change Stream Partition Metadata Engine]
↓
[Internal Log Emitter / Tokenizer]
↓
[gRPC Stream to Consumer (Dataflow / Custom Client)]

キーポイント①:トランザクションログの直接キャプチャ

よくある勘違いとして、「Change Streamsのために内部で別テーブルにトリガーを貼っている」と思う人がいるが、それは大間違いだ。そんなことをすれば、書き込み性能(TPS)は数分の一に劣化する。

実際には、Spannerのストレージエンジン(LSMツリーベースのRocksDB等)に書き込まれる先行書き込みログ(WAL)およびPaxosのコミットログのストリームを、内部のバックグラウンドプロセスが極めて低いオーバーヘッドでインターセプトしている。
これにより、データの永続化と変更ログの生成がアトミックに行われ、「書き込んだのにログが出ない」という矛盾が物理的に排除される。

キーポイント②:スプリットの動的追跡とPartitionの概念

Spannerの真骨頂は、負荷に応じてスプリットが分裂(Split)したり、統合(Merge)したりすることだ。
Change Streamsの世界では、これを Partition(パーティション) という概念で抽象化している。

  • スプリットが分裂すると、元のパーティションは「終了」し、新しい2つの子パーティションが「誕生」する。
  • Change Streamsのコンシューマ側(Dataflowなど)は、このパーティションのトポロジー変化(Directed Acyclic Graph: DAG)を常に追跡し続ける必要がある。

この「パーティションの動的追跡」を自前で実装しようとすると地獄を見る。だからこそ、Googleが提供する Cloud Dataflow テンプレート(Spanner to BigQueryなど)を使うのが実務上のデファクトスタンダードなのだ。

—

2. 堅牢な設計パターン:プロダクションで溺れないための極意

設計レビューで私が必ずチェックするポイントを授けよう。これらを外すと、数ヶ月後に本番障害の波に呑まれることになる。

パターンA:対象カラムの絞り込み(`FOR UPDATE` の罠)

Change Streamsを作る際、全てのカラムを流していなくだろうか?

— ❌ アンチパターン:テーブル全体をストリーム
CREATE CHANGE STREAM AllTablesStream FOR ALL OPTIONS (
retention_period = ‘7d’
);

— ⭕️ ベストプラクティス:必要なテーブル、必要なカラム、必要なイベントのみ
CREATE CHANGE STREAM UserAccountChanges
FOR Users(user_id, email, updated_at)
OPTIONS (
retention_period = ‘1d’
);

理由:
`FOR ALL` を指定すると、頻繁に更新される不要なメタデータカラムやログ的なカラムの変更まですべてシリアライズされ、ネットワーク帯域とコンシューマのCPUを無駄に食いつぶす。また、`retention_period`(最大7日間)を長くしすぎると、内部ストレージのコストに直結する。本当に必要な最小限の期間(通常は1日〜3日)に留めよ。

パターンB:バックプレッシャー(Backpressure)とコンシューマのスケール

Change Streamsは、サーバー側からgRPCでデータをプッシュし続ける。もし君のコンシューマ(アプリケーションやDataflowワーカー)の処理が追いつかなくなったらどうなるか?

Spanner側はバッファリングを試みるが、長期間コンシューマが遅延すると、ストリームが切断されるか、最悪の場合データロス(Retention期間超過によるパージ)が発生する。

設計上の鉄則:
1. コンシューマは必ず水平分散(Autoscale)させること。 Dataflowを使う場合は、Streaming Engineを有効化し、適切なVCPU数を割り当てる。
2. コンシューマ側での重い処理(外部APIコールや重いJOIN)をストリームのループ内で同期的に行わないこと。 一度キュー(Pub/Subなど)に突っ込んで非同期バッチ処理へ逃がせ。

—

3. パフォーマンスとコストの限界突破:プロのチューニング

最後に、パフォーマンスとコストを極限まで最適化するための知見を共有しよう。

1. 「ホットスポット」を生むスプリット分割のコントロール

Change Streamsを有効にしたテーブルに対して、単一のキー(例えばインクリメンタルなIDや、常に現在時刻を含むキー)で大量の書き込みを行うと、特定のスプリットに負荷が集中する。
結果として、そのスプリットのログを生成・配信する内部ワーカーにも負荷が偏り、Change Streamsのレイテンシ(遅延時間)が数秒〜数分に悪化する。

対策:
主キーには必ず Hash Prefix(ハッシュ化されたプレフィックス) や UUIDv4 などのランダム性の高い値を使用し、書き込みを綺麗に全スプリットへ分散させよ。スプリットが綺麗に分散していれば、Change Streamsの配信遅延はミリ秒単位(P99でも数百ミリ秒以内)に収まる。

2. データの重複と「At-Least-Once(少なくとも1回)」のハンドリング

分散システムの常として、Change Streamsの配信保証は At-Least-Once(少なくとも1回) である。ネットワークの切断やリトライにより、同じ変更イベントが複数回流れてくることがある。

コードレベルでの実装指針:
コンシューマ側でダウンストリーム(BigQueryや外部キャッシュ)へ書き込む際は、必ず べき等性(Idempotency) を担保しなければならない。

コンシューマ側(Python / Pseudocode)のべき等な処理パターンの例
def process_change_stream_event(event):
mutation_id = event.mod_type + event.keys + event.commit_timestamp

# 1. 既に処理済みのトランザクション/変更かチェック (Redis等を利用)
if cache.exists(mutation_id):
logger.info(f”Duplicate event ignored: {mutation_id}”)
return

# 2. 実際のビジネスロジック・ダウンストリームへの書き込み
downstream_db.upsert(event.data)

# 3. 処理済みフラグを保存 (TTLはChange Streamsの保持期間以上に設定)
cache.set(mutation_id, “processed”, ex=86400)

この「べき等性の担保」をサボるエンジニアが後を絶たない。コードレビューでは真っ先にここを突くこと。

—

チーフアーキテクトからの総括

Cloud Spannerの Change Streams は、単なる機能追加ではない。
「強整合性を持つグローバル分散DBの内部状態を、リアルタイムに外部へ安全にフェンスアウトする」という、分布式システムのロマンと泥臭さが詰まった最高傑作だ。

その内部構造(トランザクションログの直接キャプチャ、パーティションの動的追跡)を正しく理解し、適切なスキーマ設計(キー分散とカラム絞り込み)、そしてコンシューマ側の堅牢な実装(べき等性とバックプレッシャー対策)を行えば、あなたのシステムは圧倒的なスケーラビリティと信頼性を手に入れる。

設計に妥協するな。コードの裏側にある分散の物理法則に思いを馳せろ。君ならできる。

コメント

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