【実務・中級編】 Change Streams内部実装 – Cloud Spanner

Cloud Spanner Change Streamsの内部実装と「壊れない」分散ストリーミング設計

こんにちは。テックリードの私だ。
今日のコードレビューで、「Change Streamsからイベントを読み取って下流のKafkaやBigQueryに流すパイプライン」の設計書を見た。
……正直に言おう。甘い。あまりにもSpannerの分散アーキテクチャへの理解が不足している。

「単にテーブルの変更をフックする機能」だと思って使っているなら、今すぐその設計を止めなさい。
Cloud SpannerのChange Streamsは、単なるRDBのトランザクションログのオマケ機能ではない。これは、「数ペタバイトのスケールで整合性を完全保証しながら、分散ストレージ層から直接変更差分をブロードキャストする怪物じみた仕組み」だ。

今回は、Change Streamsの内部アーキテクチャが物理レベルでどう動いているのかを紐解き、我々エンジニアが本番環境で踏み抜く地雷を回避するための「極限の設計パターン」を伝授しよう。

—

1. 内部アーキテクチャ:なぜChange Streamsは破綻しないのか?

まず、Spannerのストレージ層の基本を思い出してほしい。
Spannerのデータは、主キーによってソートされ、Splits(スプリット)と呼ばれる物理的なシャードに分割されている。データが増えればスプリットは自動分割され、世界中の異なるリーダー・レプリカノードに分散配置される。

では、この分散された海のようなストレージから、どうやって「全テーブルの変更履歴」を漏れなく、順序を狂わせずに抽出しているのか?

物理ストレージ層(Colossus)とトランザクションログの直読

RDBの多くは、WAL(Write-Ahead Log)をシーケンシャルに舐めることで変更を検知する。しかしSpannerのストレージ基盤(Colossus)上では、データはLSMツリー(Log-Structured Merge-tree)ベースのファイル群として保持されている。

Change Streamsの裏側では、各スプリットのストレージ層レベルで以下の処理が並行稼働している。

1. Mutationのキャプチャ:
トランザクションがコミットされると、Paxosグループの合意形成を経てストレージに書き込まれる際、ミューテーション(Insert/Update/Delete)が低レベルのログとしてキャプチャされる。
2. パーテーション(Partition)の動的管理:
スプリットが分裂・統合(Split/Merge)すると、Change Streamsの論理的な読み取り単位である「パーテーション」も動的に再分割・結合される。
3. データバッファリングとWatermarking:
分散環境下では、ネットワーク遅延やノード障害により、イベントの到着順序が前後する。SpannerはWatermark(ウォーターマーク)と呼ばれるメカニズムを使い、「この時刻までの変更は完全に出揃った」と保証できる限界点を計算し、その安全なラインまでをストリームとして下流に解放する。

つまり、Change Streamsの読み取りワーカー(後述のDataflowなど)がクエリを投げているとき、彼らは単一のテーブルを見ているのではなく、「時間軸とスプリット空間で完全に調停された分散ログのストリーム」をスキャンしているのだ。

—

2. 現場で即座に使える! 堅牢なDDLとクエリ設計

アーキテクチャを理解したところで、実務の話をしよう。
Change Streamsを有効化する際、適当に `FOR ALL` なんて指定していないだろうか?

アンチパターン:全テーブル・全カラムの無差別キャプチャ

— 【悪手】これでは不要なI/Oとコストが爆発する
CREATE CHANGE STREAM AllActivityStream
FOR ALL
OPTIONS (
retention_period = ‘7d’
);

これをしていいのは、お遊びのプロトタイプ環境だけだ。本番でこれをやると、頻繁に更新されるメタデータテーブルの変更まですべてストリームに乗せられ、ストレージI/Oとクエリコスト(Compute Capacity)を無駄に食い潰すことになる。

ベストプラクティス:対象の絞り込みとRetentionの最適化

必要なテーブル、必要なカラムだけに絞り、リテンション期間もビジネス要件の最小限(通常は24時間〜最大7日)に抑えよ。

— 【推奨】対象を明示し、監査に必要なカラムだけに絞る
CREATE CHANGE STREAM UserAuditStream
FOR
Users(UserId, Email, UpdatedAt),
UserProfiles(UserId, Bio)
OPTIONS (
retention_period = ‘1d’ — データ基盤へのロード遅延を考慮しても24時間あれば十分なことが多い
);

—

3. アプリケーション設計とDataflowパイプラインの勘所

Change Streamsからデータを読み出す際、公式で推奨されているのは Google Cloud Dataflow を用いたパイプラインだ。ここで多くのエンジニアが「重複処理(At-least-once delivery)」と「順序保証」の壁にぶつかる。

なぜ「At-least-once(少なくとも1回)」なのか?

分散システムの鉄則として、Change Streamsのコネクタは「少なくとも1回」の配信を保証する。つまり、ネットワークの瞬断やワーカーのクラッシュにより、同じ変更イベントが2回流れてくる可能性が常にある。

下流システム(例えば決済システムや外部API)がこれを受け取ったとき、べき等性(Idempotency)が担保されていなければ、データが二重計上されるという大惨事を引き起こす。

設計パターン:べき等性を保証するシャドー・キーの活用

Dataflowのパイプライン内、あるいは下流のDB/DWHに書き込む際の変化検知クエリでは、必ず「一意なイベントID」を利用してべき等性を担保すること。

SpannerのChange Streamsが吐き出すレコードには、以下のメタデータが含まれている。

  • `commit_timestamp`(コミット時刻)
  • `server_transaction_id`(トランザクションID)
  • `mod_sequence_number`(同一トランザクション内のミューテーション連番)

これらを組み合わせたハッシュ値を「イベントID」として一意制約(Unique Constraint)のキーにし、下流へのUpsert(あるいはInsert Ignore)を実装する。

Dataflow (Apache Beam Python) でのべき等な処理イメージの断片
import apache_beam as beam
import hashlib

def generate_id(element):
# トランザクションIDとモッドシーケンスから一意のキーを生成
raw_str = f”{element[‘server_transaction_id’]}_{element[‘mod_sequence_number’]}”
element[‘event_id’] = hashlib.sha256(raw_str.encode(‘utf-8’)).hexdigest()
return element

with beam.Pipeline() as p:
(
p
| ‘ReadFromSpanner’ >> ReadChangeStream(…)
| ‘AddEventId’ >> beam.Map(generate_id)
| ‘IdempotentWriteToTarget’ >> beam.ParDo(IdempotentWriterDoFn())
)

この「`server_transaction_id` + `mod_sequence_number`」の組み合わせこそが、分散環境における重複排除の最強のパスポートとなる。

—

4. パフォーマンスとオペレーションの注意点

最後に、チーフアーキテクトとして、運用フェーズで絶対に踏んではいけない地雷を挙げておく。

1. キャッチアップ遅延(Lag)の監視を怠るな
下流のDataflowが何らかの原因で止まった場合、Spanner側の `retention_period` を過ぎるとデータがロストする(データ保持期間を超えた変更ログは物理的に削除されるため)。
Cloud Monitoringで `change_streams/pending_seconds`(未処理の遅延時間)を常に監視し、閾値を超えたらPagerDuty等で叩き起こされるアラート設計を義務付けろ。

2. ホットスポットを生むトランザクションに注意
1つの巨大なトランザクションで何百万行も更新すると、そのトランザクションのログを処理する単一のChange Streamパーテーションに負荷が集中し、スループットが急低下する。
バッチ処理であっても、適度なサイズ(数千〜数万行単位)に分割してコミットする原則を崩してはならない。

—

まとめ

Cloud SpannerのChange Streamsは、正しく使えば「最強のリアルタイム・データ連携基盤」になるが、その裏にある分散ストレージ、Watermark、At-least-onceのセマンティクスを理解せずに触ると、必ず本番障害という名のツケを払うことになる。

設計レビューでこのあたりの整合性やリトライ・べき等性の担保が抜けているコードを見かけたら、私の名前を出して差し戻してくれて構わない。

妥協のない、堅牢な分散システムを作ろう。健闘を祈る。

コメント

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