PROD欧州主権の BaaS プラットフォームダッシュボードを開く →

エンジニアリング · 10 分読み取り

PostgreSQL CDC と NATS: リアルタイム アーキテクチャ

Affane Daylami · Fondateur · 2026年7月24日

ブログに戻る

WebSocket 経由で接続された Aurabase クライアントでは、ポーリングや Webhook の設定を行わずに、COMMIT の数ミリ秒後にベースに挿入された行が表示されます。仕組みは3段階になっています。論理レプリケーション プラグインは Postgres WAL を JSON にデコードし、1 つのプロセスのみがそれを読み取る権利を持ち、NATS が各イベントをクラスター内のすべての WebSocket サーバーに中継します。 aura-realtimeservice の実際のコードでどのように接続されているかを次に示します。

この英語のテキストはフランス語のオリジナルから自動的に生成されたもので、まだレビューされていません。
このページは自動翻訳されました。英語版が正式です。

ここにあるものは理想的なパターンではありません。各抽出物は、発行日の aura-realtimeクレートからそのまま提供されます。ここには明確な違いがあります。NATS JetStream は、このパイプラインで CDC イベントを転送するものではありません。以下では、JetStream が実際に何を行うのか、また代わりに CDC トラフィックを伝送するものについて詳しく説明します。エラーを発生させずにメカニズム全体をサイレントにすることができる検証済みの Kubernetes トラップについては、以下で詳しく説明します。

必需品
  • Aurabase の Postgres CDC は、pgoutputバイナリ プロトコルや Debezium/Kafka Connect ブリッジではなく、wal2json および pg_logical_slot_get_changes 経由で WAL を読み取ります。
  • 論理スロット aura_cdc_slot は 1 つのリーダーのみを許容します。 aura-realtime-cdc-worker は、K8 を除く「常にリーダー」モードで、Kubernetes リース (coordination.k8s.io/v1) 経由で選出されます。
  • ws-front レプリカへのファンアウトは、永続的な JetStream ストリームではなく、NATS コア ( aura.realtime.>上の単純なパブリッシュ/サブスクライブ) を経由します。 JetStream は、この同じサービス内でクロスインスタンス プレゼンス KV を排他的に提供します。
  • 各イベントは、キャッシュされた近似値ではなく、該当する回線上の実際の SET LOCAL ROLE リクエストを介して、送信直前にサブスクライバによって RLS で再チェックされます。
  • リポジトリ履歴の落とし穴を確認しました。Postgres Kubernetes イメージに wal2json プラグインがないと、スロットの作成が失敗し、CDC が静かに非アクティブなままになります。
#
エンジンの選択

wal2json、pgoutput、Debezium ではありません

ほとんどの CDC Postgres パイプラインは、バイナリ ロジック レプリケーション プロトコルである pgoutputを通過し、それを Kafka に変換する Debezium などのコネクタを通過します。 Aurabase はこのステップを省略します。 aura-realtime サービスは、 、 wal2json プラグインを使用して、 WAL を論理レプリケーション に直接デコードし、中間の変換段階なしで使用可能な JSON を生成します。

cdc/postgres.rs (実際のクエリ)sql
-- 存在しない場合は一度作成され、WAL の有効期限が切れた場合は再作成されます
SELECT pg_create_logical_replication_slot(
  'aura_cdc_slot', 'wal2json');

-- ポーリングを繰り返します (バッチが空になったときに 100 ミリ秒バックオフします)
SELECT data FROM pg_logical_slot_get_changes(
  'aura_cdc_slot', NULL, 1000,
  'include-xids', '0', 'format-version', '2',
  'include-schemas', '1', 'include-types', '0');

コードで確認された 3 つの Postgres 前提条件: アプリケーション ロールは REPLICATIONを保持する必要があり、 wal_level = logical がアクティブである必要があり、 max_slot_wal_keep_size が WAL 保持を制限する必要があります。この最後の制限がないと、コンシューマーが遅いため、ディスク ボリュームが無制限に増加します。このサービスはスロットの wal_status も監視します。スロットが lost (この制限を超えて WAL がパージされる) に変化すると、スロットは自動的に再作成されます。その間消費されなかったイベントの損失が想定され、ログに記録されます。

動的テーブル フィルターが 3 秒ごとに更新される、ポーリングごとに 1000 件の変更のバッチ。プロジェクトで明示的にリアルタイムが有効になっているテーブルのみ add-tables を入力します。これにより、サブスクライバーにとって関係のないテーブルの WAL のデコードが回避されます。

#
一人のプロデューサー

cdc-worker: Kubernetes リースによって選出され、重複していません

Postgres 論理レプリケーション スロットは、一度に 1 つのアクティブ ドライブのみを許容します。同じ aura_cdc_slot の複数のコンシューマーを並行して実行すると、変更が重複するだけでなく、変更の順序が崩れてしまいます。 Aurabase は、WebSocket の水平スケールを犠牲にすることなくこの制約を解決するために、aura-realtime を 2 つの個別のバイナリに分割しました。

Cargo.tomltoml
# cdc-worker: リーダーに選出され、CDC を立ち上げ、NATS で公開します。 WSサーバーがありません。
[[bin]]
name = "aura-realtime-cdc-worker"

# ws-front: N レプリカ (HPA)、NATS を消費、RLS フィルター、WS/SSE を提供します。
[[bin]]
name = "aura-realtime-ws-front"

cdc-worker は、リース Kubernetes オブジェクト (coordination.k8s.io/v1) を保持することによってのみ、スロットを読み取る権利を取得します。彼はそれを作成しようとするか、有効期限が切れている場合は盗もうとし、寿命の 3 分の 1 が経過すると更新します。リースが失われると (更新が失敗し、所有者が変更されると)、CancellationToken によって進行中の作業が中断され、プロセスは取得ループに戻ります。 Kubernetes の外部では (ローカル テスト中や非クラスター開発中など)、K8s クライアントは接続に失敗し、サービスは「常にリーディング」モードに切り替わります。開発環境では便利ですが、複数のレプリカがある本番環境で忘れた場合は正しくありません。

cdc/leader.rs (本物の署名)rust
pub async fn run_with_lease<F, Fut>(cfg: LeaseConfig, work: F) -> anyhow::Result<()>
where F: Fn(CancellationToken) -> Fut, ...
{
  // すべての release_duration_secs / 3 を更新します
  // K8 を除く: ワーク (トークン) は 1 回のみ、トークンはキャンセルされません
}
#
リアルファンアウト

CDC 用の NATS コア、プレゼンス用の JetStream

これは、この種のパイプラインで最も誤解されやすい点です。NATS と JetStream は 2 つの別個のものであり、CDC イベントは永続的な JetStream ストリームを通過しません。 cdc-worker は、Client::publish_with_headers — NATS コア パブ/サブ APIを使用して各イベントをパブリッシュします。これは、JetStream とは別に、永続化やリプレイなしで 最大 1 回 で配信されます。

NATS ハート (パブ/サブ)CDC イベントの ws-front レプリカへのファンアウト永続化や再実行なしで、最大 1 回のみ
NATS ジェットストリーム (KV)クロスインスタンスプレゼンス (aura_presence)共有状態 60 秒、再生する CDC ストリームではない
channels/manager.rs → nats/mod.rs (extraits réels)rust
// cdc-worker がパブリッシュ (ChannelManager::publish_and_relay)
nc.publish_with_headers(subject, headers, payload).await

// ws-front はサブツリー全体 (NatsRelay::start_relay) をリッスンします。
client.subscribe("aura.realtime.>").await

各レプリカ ws-front はサブツリー aura.realtime.>をサブスクライブし、トピックを元のチャネル名にデコードします。次に、接続されている WebSocket のローカル tokio::sync::broadcastにイベントを再発行します。 Aura-Origin ヘッダーには、送信インスタンスの識別子が含まれます。各レプリカは、自身が発行したメッセージを無視します。これにより、中央調整なしで重複が回避されます。

この選択は直接的な結果をもたらします。NATS コアは、長時間持続するキューを持たずに、最大 1 回 (最大 1 回) を配信します。イベントの通過時に ws-front レプリカが NATS から一時的に切断されると、追いつきません。チャネルごとの履歴 (Channel.history、 HISTORY_SIZEによって境界付けられたメモリ内バッファ) は、クラスター レベルではなく、各レプリカでローカルに存在します。 Postgres 論理レプリケーション スロットが上流の永続的な信頼できるソースであり続けるため、これは許容できる妥協策です。使用前に DB の変更が失われないようにするのは、NATS ではなく wal2json です。

JetStream は aura-realtime に存在しますが、用途はまったく異なります。これは、クロスインスタンス プレゼンス KV ストア (aura_presence、メモリ ストア、60 秒 max_age) を提供し、すべての ws-frontレプリカ間で誰がどのチャネルでオンラインであるかを同期します。短期間の共有状態であり、再生する CDC イベントのストリームではありません。

#
セキュリティ

イベントごとに加入者によって RLS が再チェックされる

CDC イベントは、チャネルのすべての加入者にそのまま送信されるわけではありません。ポーラーは各変更を CDC_NUM_SHARDS ワーカー (デフォルトでは 4 つ、 project_idによってハッシュ化) の 1 つにルーティングし、最初に pg_policiesをポーリングします。 RLS ポリシーのないテーブルではすべてのサブスクライバー (Supabase の動作) が許可されますが、ポリシーのあるテーブルではサブスクライバーごとのチェックがトリガーされます。

subscriptions.rs::check_rls_visible_batch (実際の、簡略化された)sql
BEGIN;
SET LOCAL ROLE "aura_authenticated";
-- GUC に注入されたサブスクライバの JWT クレーム
SELECT set_config('request.jwt.claims', $1, true);
SELECT EXISTS(SELECT 1 FROM "schema"."table" WHERE id::text = $2);
ROLLBACK;

トランザクションはサブスクライバーの物理的な Postgres の役割を引き受け、その JWT クレームを request.jwt.claims に挿入します。これは、auth.uid() および auth.role() がポリシー側で読み取るものです。次に、このロールの下でラインが表示されたままであることを確認し、書き込みを行わず、実際の RLS 読み取りをキャンセルします。検証された sub_id のみがイベントの allowed_sub_ids に含まれます。空の vec は、フェールクローズでは誰もいないことを意味します。

アストゥチェ

DELETE は、このチェックを行ごとにエスケープし、行が消えているため、RLS で再読み取りすることは不可能であり、チャネルのすべての加入者に戻ります。その代わり、DELETE のペイロードは ID 列 (主キー) のみを公開し、削除された行の内容は公開しません。

移行の邪魔になることが多い詳細: wal2json は、テーブルに REPLICA IDENTITY FULLがある場合、 UPDATE または DELETE の完全な内容のみをキャプチャします。これがないと、INSERT のみが復帰します。これが、テーブルでリアルタイムを有効にするエンドポイントがこの ALTER TABLE を自動的に設定する理由です。これは、個別のチェックボックスではなく、リポジトリに対する実際の修正です。

#
トラップが確認されました

CDC が Kubernetes で沈黙を守ることができる理由

この出願履歴は、理論上の事件ではなく、実際の罠を記録している。 Kubernetes クラスターでデフォルトで使用される Postgres イメージ (pgvector/pgvector ストック コミュニティ イメージ) には、wal2jsonプラグインが含まれていません。これがないと、pg_create_logical_replication_slot('aura_cdc_slot', 'wal2json') は失敗し、クライアント側ではこれを明示的に通知するものはありません。 WebSocket は開いたままで、サブスクリプションは受け入れられますが、DB イベントは発生しません。

docker/Postgres.Dockerfile (本物のエキス)dockerfile
RUN apt-get update && \
    apt-get install -y --no-install-recommends \
      postgresql-16-wal2json=2.6-4.pgdg12+1 \
      ...

実際の修正は、このパッケージがインストールされた専用の Postgres イメージを構築して公開し、ストック イメージの代わりに Kubernetes StatefulSet をそのイメージに指定することでした。また、汎用イメージには存在しない、 wal_level=logical、 max_replication_slots 、および max_slot_wal_keep_size を起動引数として明示的に設定します。

cdc-worker が専用の /health エンドポイントで cdc_active フラグ (スロットが正常であることが確認されない限り false) を公開するのは、まさにこのタイプの障害を観察可能にするためです。これは、クライアント側イベントの欠如だけから CDC が停止していると推測するのではなく、注意すべきバイナリ シグナルです。

#
電流制限

このパイプラインでまだカバーされていないもの: 専用クラスター

このメカニズムは単一の POSTGRES_REPLICATION_URL変数を読み取ります。つまり、単一の Postgres サーバーと単一の論理スロットが読み取ります。これは、共有クラスター上のプロジェクトごとのマルチテナント スキーマ モデル (project_<uuid>) と一致しています。このモデルでは、すべてが機能します。単一の cdc-worker は、同じクラスター内のすべてのプロジェクトの変更を確認し、図ごとにルートします。

専用の CNPG クラスター (独自の分離された Postgres インスタンス) 上のプロジェクトの場合、現在使用されているイメージ (CloudNativePG ストック イメージから構築) は wal2jsonではなく pg_graphqlのみを追加します。また、専用クラスターごとに cdc-worker を実行するものはありません。したがって、コードの現在の状態では、ここで説明するリアルタイム CDC は、ロードマップ機能ではなく、実際のアーキテクチャ境界である共有層に焦点を当てたままになります。

#
クライアント側

変わらないもの: SDK の postgres_changes API

このメカニズム全体は SDK からは見えないままです。aurabase-js チェーン可能 API は、イベントが NATS 経由の cdc-worker から来たのか、ローカル開発で分割されていない aura-realtime バイナリから来たのかに応じて変わりません。

app.tstypescript
const channel = aura.realtime.channel('room-1')
  .on('postgres_changes', {
    event: 'INSERT', schema: 'public', table: 'messages'
  }, (payload) => { /* … */ })
  .subscribe()

ただし、単一のイベントが発生する前に、テーブルを CDC に登録する必要があります。これは作成時に自動的に行われるわけではありません。専用エンドポイント (または Studio の同等のトグル) での service_role 呼び出しは、プロジェクトのプラットフォーム スキーマの realtime_tables にテーブルを書き込み、REPLICA IDENTITY FULL を設定します。

terminalbash
curl -X PUT \
  https://api.aurabase.cloud/v1/db/<project_id>/schema/tables/messages/realtime \
  -H "apikey: <service_role_key>" \
  -H 'Content-Type: application/json' \
  -d '{"enabled": true}'

Supabase と互換性のある SDK サーフェスの残りの部分 (認証、ストレージ、ポリシー RLS) については、 移行ガイドを参照してください。 aura-realtime が単一のワークスペース クレート内で複数のバイナリに分割される方法については、 Cargo アーキテクチャの詳細を参照してください。

導入の準備はできていますか?

5 分でバックエンドが完成します。

クレジット カードは不要 · 500 MB 無料 · 50,000 MAU