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のセマンティクスを理解せずに触ると、必ず本番障害という名のツケを払うことになる。
設計レビューでこのあたりの整合性やリトライ・べき等性の担保が抜けているコードを見かけたら、私の名前を出して差し戻してくれて構わない。
妥協のない、堅牢な分散システムを作ろう。健闘を祈る。
コメント