•18 min read

PostgreSQLの変更データキャプチャ(CDC): Debezium、Kafka Connect、トランザクションアウトボックス

PostgreSQLの変更データキャプチャ(CDC): Debezium、Kafka Connect、トランザクションアウトボックス

このガイドでは、Debezium、Kafka Connect、およびトランザクショナルアウトボックスパターンを活用した、PostgreSQL向け堅牢な変更データキャプチャ(CDC)パイプラインの実装について詳しく説明します。目的は、二重書き込みの不整合をゼロにし、信頼性の高いデータ伝播を保証するイベント駆動型マイクロサービスアーキテクチャを構築することです。

Audio Briefing
0:00 / 0:00

論理レプリケーションのためのPostgreSQL設定

DebeziumはPostgreSQLの論理デコーディング機能に依存しています。これには特定のサーバー設定が必要です。

postgresql.confの調整

論理デコーディングを有効にするには、wal_levelパラメータをlogicalに設定する必要があります。DebeziumのレプリケーションスロットとWALセンダープロセスに対応するために、max_replication_slotsとmax_wal_sendersを設定する必要があります。

# postgresql.conf
wal_level = logical
max_replication_slots = 10 # Adjust based on number of Debezium connectors
max_wal_senders = 10       # Adjust based on number of Debezium connectors

postgresql.confを変更した後、これらの変更を有効にするにはPostgreSQLの再起動が必須です。

レプリケーションスロットの作成

Debeziumは変更を追跡するために論理レプリケーションスロットを必要とします。このスロットは、Debeziumが処理するまでPostgreSQLが必要なWrite-Ahead Log(WAL)セグメントを保持することを保証し、データ損失を防ぎます。pgoutputプラグインは論理デコーディングの標準です。

-- Connect as a superuser or a user with REPLICATION privileges
SELECT * FROM pg_create_logical_replication_slot('debezium_slot', 'pgoutput');

スロットの存在と状態を確認します。

SELECT slot_name, plugin, active, wal_status, restart_lsn FROM pg_replication_slots;

restart_lsnは、スロットが変更のストリーミングを開始するWALの位置を示します。wal_statusは理想的にはreservedまたはactiveであるべきです。activeがfの場合、Debeziumは現在接続されていません。

Advertisement

Debezium Kafka Connectのセットアップ

Debeziumは、データベースをイベントストリームに変換する分散プラットフォームです。Kafka Connectと統合して、PostgreSQLからの変更をKafkaトピックにストリーミングします。

Kafka Connectのデプロイ

Kafka Connectは、スタンドアロンモードまたは分散モードでデプロイできます。本番環境では、フォールトトレランスとスケーラビリティのために分散モードが推奨されます。

# Example: Distributed Kafka Connect worker configuration (connect-distributed.properties)
bootstrap.servers=kafka-broker-1:9092,kafka-broker-2:9092
group.id=connect-cluster
config.storage.topic=connect-configs
offset.storage.topic=connect-offsets
status.storage.topic=connect-statuses
key.converter=io.confluent.connect.avro.AvroConverter
key.converter.schema.registry.url=http://schema-registry:8081
value.converter=io.confluent.connect.avro.AvroConverter
value.converter.schema.registry.url=http://schema-registry:8081
internal.key.converter=org.apache.kafka.connect.json.JsonConverter
internal.value.converter=org.apache.kafka.connect.json.JsonConverter
internal.key.converter.schemas.enable=false
internal.value.converter.schemas.enable=false
plugin.path=/usr/share/java,/usr/share/confluent-hub-components

Debezium PostgreSQLコネクタのJARがplugin.pathで指定されたディレクトリに配置されていることを確認してください。

Debezium PostgreSQLコネクタの設定

Debeziumコネクタの設定は、PostgreSQLへの接続方法とストリーミングするデータを定義します。

{
  "name": "postgres-connector",
  "config": {
    "connector.class": "io.debezium.connector.postgresql.PostgresConnector",
    "tasks.max": "1",
    "database.hostname": "postgres-db",
    "database.port": "5432",
    "database.user": "debezium_user",
    "database.password": "debezium_password",
    "database.dbname": "your_database",
    "database.server.name": "your_logical_server_name",
    "schema.include.list": "public",
    "table.include.list": "public.users,public.orders",
    "slot.name": "debezium_slot",
    "publication.name": "debezium_publication",
    "publication.autocreate.mode": "all_tables",
    "plugin.name": "pgoutput",
    "topic.prefix": "dbserver",
    "heartbeat.interval.ms": "5000",
    "snapshot.mode": "initial",
    "decimal.handling.mode": "double",
    "time.precision.mode": "connect",
    "hstore.handling.mode": "json",
    "converters": "dateConverter",
    "dateConverter.type": "org.apache.kafka.connect.data.Timestamp",
    "dateConverter.format": "yyyy-MM-dd HH:mm:ss.SSS",
    "key.converter": "io.confluent.connect.avro.AvroConverter",
    "key.converter.schema.registry.url": "http://schema-registry:8081",
    "value.converter": "io.confluent.connect.avro.AvroConverter",
    "value.converter.schema.registry.url": "http://schema-registry:8081"
  }
}

この設定をKafka Connect REST API経由でデプロイします。

curl -X POST -H "Content-Type: application/json" --data @debezium-connector.json http://kafka-connect:8083/connectors

Confluent Schema RegistryとAvroによるスキーマ進化

Avroコンバータで設定されたDebeziumは、Confluent Schema Registryにスキーマを自動的に登録します。これは、イベントストリームにおけるスキーマ進化を管理するために非常に重要です。テーブルスキーマが変更された場合(例:列の追加)、Debeziumはこれを検出し、Schema RegistryのAvroスキーマを更新し、後続のメッセージは新しいスキーマを反映します。これにより、コンシューマはSchema Registryを使用して、スキーマバージョンが異なっていてもメッセージを正しく逆シリアル化できます。

トランザクショナルアウトボックスパターン

トランザクショナルアウトボックスパターンは、二重書き込みの問題に対処します。つまり、データベーストランザクションと送信メッセージの公開がアトミックであることを保証します。Kafkaに直接公開する代わりに、イベントはビジネスロジックと同じデータベーストランザクション内でoutboxテーブルに書き込まれます。その後、Debeziumがこれらのoutboxテーブルの変更を検出し、Kafkaに公開します。

アウトボックステーブルスキーマ

CREATE TABLE outbox (
    id UUID PRIMARY KEY,
    aggregatetype VARCHAR(255) NOT NULL,
    aggregateid UUID NOT NULL,
    type VARCHAR(255) NOT NULL,
    payload JSONB NOT NULL,
    createdat TIMESTAMP WITH TIME ZONE DEFAULT CURRENT_TIMESTAMP
);

パターンの実装

新しいユーザーを作成し、UserCreatedイベントを公開する必要があるUserServiceを考えてみましょう。

import { Pool } from 'pg';
import { v4 as uuidv4 } from 'uuid';

interface User {
  id: string;
  name: string;
  email: string;
}

interface OutboxEvent {
  id: string;
  aggregateType: string;
  aggregateId: string;
  type: string;
  payload: any;
}

class UserService {
  private pool: Pool;

  constructor(pool: Pool) {
    this.pool = pool;
  }

  public async createUser(name: string, email: string): Promise<User> {
    const client = await this.pool.connect();
    try {
      await client.query('BEGIN');

      const userId = uuidv4();
      const user: User = { id: userId, name, email };

      // 1. Persist business entity
      await client.query(
        'INSERT INTO users (id, name, email) VALUES ($1, $2, $3)',
        [user.id, user.name, user.email]
      );

      // 2. Create outbox event
      const event: OutboxEvent = {
        id: uuidv4(),
        aggregateType: 'User',
        aggregateId: user.id,
        type: 'UserCreated',
        payload: { userId: user.id, name: user.name, email: user.email, timestamp: new Date().toISOString() },
      };

      // 3. Persist outbox event in the same transaction
      await client.query(
        'INSERT INTO outbox (id, aggregatetype, aggregateid, type, payload) VALUES ($1, $2, $3, $4, $5)',
        [event.id, event.aggregateType, event.aggregateId, event.type, JSON.stringify(event.payload)]
      );

      await client.query('COMMIT');
      return user;
    } catch (error) {
      await client.query('ROLLBACK');
      console.error('Failed to create user and publish event:', error);
      throw error;
    } finally {
      client.release();
    }
  }
}

// Example usage (assuming 'pool' is a configured pg.Pool instance)
// const userService = new UserService(pool);
// userService.createUser('John Doe', 'john.doe@example.com')
//   .then(user => console.log('User created:', user))
//   .catch(err => console.error('Error:', err));

outboxテーブルを監視するように設定されたDebeziumは、INSERT操作をキャプチャし、Kafkaメッセージとして公開します。

イベントコンシューマのための冪等性キー

Kafkaからイベントを消費する場合、イベントを複数回処理しても(再試行やコンシューマのリバランスによる)一貫性のない状態にならないようにすることが重要です。冪等性キーがこれを解決します。

outboxテーブルのidフィールドは、自然な冪等性キーとして機能します。コンシューマは処理されたイベントのidを保存し、重複を拒否する必要があります。

import { Kafka, Consumer } from 'kafkajs';
import { Pool } from 'pg';

interface UserCreatedEvent {
  userId: string;
  name: string;
  email: string;
  timestamp: string;
}

class UserEventHandler {
  private consumer: Consumer;
  private pool: Pool;

  constructor(kafka: Kafka, pool: Pool) {
    this.consumer = kafka.consumer({ groupId: 'user-service-consumer-group' });
    this.pool = pool;
  }

  public async start(): Promise<void> {
    await this.consumer.connect();
    await this.consumer.subscribe({ topic: 'dbserver.public.outbox', fromBeginning: false });

    await this.consumer.run({
      eachMessage: async ({ topic, partition, message }) => {
        if (!message.value) return;

        const event = JSON.parse(message.value.toString());
        const outboxEventId = event.payload.id; // The ID from the outbox table, our idempotency key
        const operation = event.op; // 'c' for create, 'u' for update, 'd' for delete

        if (operation !== 'c') {
          // We only care about new outbox events, not updates/deletes to the outbox table itself
          return;
        }

        const payload = event.after.payload; // Debezium 'after' field contains the new row
        const eventType = event.after.type;

        if (eventType === 'UserCreated') {
          const userCreatedEvent: UserCreatedEvent = payload;
          await this.processUserCreatedEvent(outboxEventId, userCreatedEvent);
        }
        // Handle other event types
      },
    });
  }

  private async processUserCreatedEvent(outboxEventId: string, event: UserCreatedEvent): Promise<void> {
    const client = await this.pool.connect();
    try {
      await client.query('BEGIN');

      // Check for idempotency: has this event ID been processed before?
      const checkResult = await client.query(
        'SELECT 1 FROM processed_events WHERE event_id = $1',
        [outboxEventId]
      );

      if (checkResult.rows.length > 0) {
        console.log(`Event ${outboxEventId} already processed. Skipping.`);
        await client.query('ROLLBACK'); // Rollback the transaction as nothing new was done
        return;
      }

      // Process the event (e.g., create a user in a read model, send a welcome email)
      console.log(`Processing UserCreated event for user ${event.userId}:`, event);
      // Example: Insert into a read-model table
      await client.query(
        'INSERT INTO read_model_users (id, name, email) VALUES ($1, $2, $3) ON CONFLICT (id) DO NOTHING',
        [event.userId, event.name, event.email]
      );

      // Record the event ID as processed
      await client.query(
        'INSERT INTO processed_events (event_id, processed_at) VALUES ($1, NOW())',
        [outboxEventId]
      );

      await client.query('COMMIT');
      console.log(`Successfully processed event ${outboxEventId}`);
    } catch (error) {
      await client.query('ROLLBACK');
      console.error(`Error processing event ${outboxEventId}:`, error);
      throw error;
    } finally {
      client.release();
    }
  }
}

// Ensure 'processed_events' table exists for idempotency tracking
// CREATE TABLE processed_events (
//     event_id UUID PRIMARY KEY,
//     processed_at TIMESTAMP WITH TIME ZONE DEFAULT CURRENT_TIMESTAMP
// );

// Example usage
// const kafka = new Kafka({ brokers: ['kafka-broker-1:9092'] });
// const userEventHandler = new UserEventHandler(kafka, pool);
// userEventHandler.start().catch(console.error);

アーキテクチャ比較:CDC vs. ポーリング vs. 二重書き込み

機能CDC (Debezium + Outbox)ポーリング二重書き込み (直接Kafka)
アトミック性保証される (単一DBトランザクション)N/A (DBクエリとメッセージ送信は別々)なし (DB書き込みとKafka送信は別々の操作)
レイテンシほぼリアルタイム定期的 (数分から数時間)ほぼリアルタイム
リソース使用量DB上は低 (WALベース)、Connect上は中DB上は高 (繰り返しクエリ)DB上は低、アプリケーションサービス上は中
データ損失リスク最小限 (WAL + レプリケーションスロット)高 (ポーリング間の変更を見逃す)高 (DB書き込みとKafka送信の間にアプリがクラッシュした場合)
複雑性中 (Debezium, Kafka Connect, Schema Registry)低 (シンプルなクエリ)低〜中 (Kafkaクライアント統合)
スキーマ進化非常に優れている (Avro + Schema Registry)手動 (慎重なクエリ/コンシューマの更新が必要)手動 (慎重なプロデューサ/コンシューマの更新が必要)
スケーラビリティ高 (Kafka Connect分散)制限あり (ポーリング頻度でDB負荷が増加)高 (Kafkaはうまくスケールする)
ユースケースイベント駆動型マイクロサービス、データウェアハウジング、監査シンプルで低ボリュームの統合重要なデータの一貫性には非推奨
Advertisement

本番環境での落とし穴とトラブルシューティング

  1. レプリケーションスロットの肥大化:

    • 障害モード: pg_replication_slotsが長期間wal_statusとしてunreservedまたはretainedを示し、restart_lsnが進まず、WALセグメントのためにディスク容量が満杯になる。
    • 原因: Debeziumコネクタがダウン、一時停止、またはスタックしており、レプリケーションスロットからの消費に失敗している。PostgreSQLは非アクティブなスロットのためにWALを無期限に保持する。
    • 修正:
      • Debeziumコネクタを再起動する。
      • Kafka Connectのログでエラー(例:ネットワークの問題、Kafkaブローカーの利用不可、Schema Registryの問題)を確認する。
      • コネクタがすぐに回復できない場合、レプリケーションスロットを削除し(最終手段として、スロットが停止してからの変更は失われる)、再作成することを検討する。SELECT pg_drop_replication_slot('debezium_slot');
      • スロットの遅延についてpg_replication_slotsとpg_wal_lsn_diff(pg_current_wal_lsn(), restart_lsn)を監視する。
  2. Schema Registryの非互換性:

    • 障害モード: Debeziumコネクタが起動に失敗するか、Kafkaで「Schema not found」や「Incompatible schema」のようなエラーを伴う読み取り不能なメッセージを生成する。
    • 原因: Schema Registryが利用できないか、互換性のないスキーマ変更(例:適切なAvro互換性設定なしに列の型をVARCHARからINTEGERに変更する)がある。
    • 修正:
      • Schema Registryの可用性とネットワーク接続を確認する。
      • Schema Registryのログを確認する。
      • コネクタ設定でvalue.converter.schema.registry.urlが正しいことを確認する。
      • 破壊的変更の場合、新しいトピック、新しいコネクタ、またはスキーマ進化を処理するための戦略(例:Avroシリアル化の前にデータを変換するためのカスタムSMTを使用する)を検討する。Confluent Schema Registryの互換性モード(例:BACKWARD、FORWARD、FULL)がここで重要になる。
  3. Kafka Connectワーカーの障害:

    • 障害モード: コネクタが処理を停止し、タスクがFAILEDとマークされるか、Kafka Connectクラスター全体が応答しなくなる。
    • 原因: メモリ不足エラー、ネットワークパーティション、設定ミスのあるプラグイン、または基盤となるKafka/ZooKeeperの問題。
    • 修正:
      • Kafka ConnectワーカーのログでOOMエラーやその他の例外を監視する。必要に応じてJVMヒープサイズを調整する。
      • KafkaおよびZooKeeperクラスターの健全性を確認する。
      • offset.storage.topic、config.storage.topic、status.storage.topicが健全でアクセス可能であることを確認する。
      • 分散モードでは、Kafka Connectはタスクを自動的にリバランスするはずですが、永続的な障害はより深い問題を示しています。
  4. アウトボックステーブルの処理遅延:

    • 障害モード: イベントがoutboxテーブルに書き込まれるが、Kafkaトピックにすぐには表示されない。
    • 原因: Debeziumコネクタが遅い、一時停止している、またはKafkaからのバックプレッシャーを受けている。outboxテーブルが過度に増加している可能性がある。
    • 修正:
      • Debeziumコネクタのステータスとログを確認する。
      • Kafkaブローカーの健全性とDebeziumトピックのコンシューマグループの遅延を監視する。
      • outboxテーブルがtable.include.listに含まれており、publication.autocreate.modeが正しく設定されていることを確認する。
      • 古いイベントを削除する予定がある場合(ただし、Debeziumはその操作にこれを必要としない)、outboxテーブルのcreatedatにインデックスを追加することを検討する。

よくある質問

  1. ビジネステーブルを直接監視する代わりに、なぜトランザクショナルアウトボックスパターンを使用するのですか? ビジネステーブルを直接監視することは単純なケースでは機能しますが、トランザクショナルアウトボックスパターンは専用の、疎結合なイベントストリームを提供します。これにより、明示的なイベント契約を定義し、追加のメタデータでイベントをエンリッチし、公開するイベントペイロードが単なる生のデータベース行の変更ではなく、意図したとおりのものであることを保証できます。また、一貫したイベント構造を提供することで、コンシューマロジックを簡素化します。

  2. DebeziumはPostgreSQLテーブルのスキーマ変更をどのように処理しますか? スキーマが変更された場合(例:列の追加、型の変更)、Debeziumはこれを検出します。AvroコンバータとSchema Registryで設定されている場合、そのトピックのAvroスキーマの新しいバージョンを登録します。その後、コンシューマはSchema Registryを使用して各メッセージの正しいスキーマバージョンを取得し、Schema Registryの互換性設定に応じて後方互換性と前方互換性を可能にします。

  3. Debeziumコネクタがダウンした場合、データは失われますか? いいえ、データは失われません。PostgreSQLの論理レプリケーションスロットは、Debeziumが正常に処理するまでWALセグメントが保持されることを保証します。Debeziumが再起動すると、オフセットストレージの最後にコミットされたLSN(Log Sequence Number)から再開し、中断したところから正確に処理を続行します。ただし、長時間のダウンタイムはPostgreSQLサーバー上のWALの肥大化につながる可能性があります。

  4. outboxテーブルから古いイベントをクリーンアップするにはどうすればよいですか? Debeziumはoutboxテーブルから読み取るだけで、レコードを削除することはありません。outboxテーブルから古いイベントを定期的に削除するには、別のプロセス(例:スケジュールされたジョブ)が必要です。このジョブは、十分に古く、Debeziumによって処理されたことが確認されたイベント(例:Debeziumのオフセットまたはタイムスタンプを確認する)のみを削除する必要があります。Debeziumがまだストリーミングしていないイベントを削除しないように注意してください。一般的な戦略は、特定のしきい値(例:7日)よりも古いイベントを削除することです。

  5. pgoutput以外の論理デコーディングプラグインでもDebeziumを使用できますか? はい、Debeziumはwal2jsonなどの他のプラグインもサポートしています。ただし、pgoutputはPostgreSQL 10で導入されたネイティブの論理デコーディングプラグインであり、その効率性とPostgreSQLのレプリケーションプロトコルとの密接な統合から、一般的に推奨されています。wal2jsonはより人間が読みやすいJSON出力を提供し、デバッグや特定のユースケースに役立つかもしれませんが、pgoutputはDebeziumを使用した本番CDCの標準です。

Share this article:

Stay Updated

Get the latest posts delivered straight to your inbox.

Free Developer Utilities

Free In-Browser Developer Tools

Clean AI CLI logs, build cron expressions, decode JWTs, and calculate chmod permissions offline.

Explore Tools
Advertisement