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

目次(14 項目)
このガイドでは、Debezium、Kafka Connect、およびトランザクショナルアウトボックスパターンを活用した、PostgreSQL向け堅牢な変更データキャプチャ(CDC)パイプラインの実装について詳しく説明します。目的は、二重書き込みの不整合をゼロにし、信頼性の高いデータ伝播を保証するイベント駆動型マイクロサービスアーキテクチャを構築することです。
論理レプリケーションのための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は現在接続されていません。
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はうまくスケールする) |
| ユースケース | イベント駆動型マイクロサービス、データウェアハウジング、監査 | シンプルで低ボリュームの統合 | 重要なデータの一貫性には非推奨 |
本番環境での落とし穴とトラブルシューティング
-
レプリケーションスロットの肥大化:
- 障害モード:
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)を監視する。
- 障害モード:
-
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)がここで重要になる。
-
Kafka Connectワーカーの障害:
- 障害モード: コネクタが処理を停止し、タスクが
FAILEDとマークされるか、Kafka Connectクラスター全体が応答しなくなる。 - 原因: メモリ不足エラー、ネットワークパーティション、設定ミスのあるプラグイン、または基盤となるKafka/ZooKeeperの問題。
- 修正:
- Kafka ConnectワーカーのログでOOMエラーやその他の例外を監視する。必要に応じてJVMヒープサイズを調整する。
- KafkaおよびZooKeeperクラスターの健全性を確認する。
offset.storage.topic、config.storage.topic、status.storage.topicが健全でアクセス可能であることを確認する。- 分散モードでは、Kafka Connectはタスクを自動的にリバランスするはずですが、永続的な障害はより深い問題を示しています。
- 障害モード: コネクタが処理を停止し、タスクが
-
アウトボックステーブルの処理遅延:
- 障害モード: イベントが
outboxテーブルに書き込まれるが、Kafkaトピックにすぐには表示されない。 - 原因: Debeziumコネクタが遅い、一時停止している、またはKafkaからのバックプレッシャーを受けている。
outboxテーブルが過度に増加している可能性がある。 - 修正:
- Debeziumコネクタのステータスとログを確認する。
- Kafkaブローカーの健全性とDebeziumトピックのコンシューマグループの遅延を監視する。
outboxテーブルがtable.include.listに含まれており、publication.autocreate.modeが正しく設定されていることを確認する。- 古いイベントを削除する予定がある場合(ただし、Debeziumはその操作にこれを必要としない)、
outboxテーブルのcreatedatにインデックスを追加することを検討する。
- 障害モード: イベントが
よくある質問
-
ビジネステーブルを直接監視する代わりに、なぜトランザクショナルアウトボックスパターンを使用するのですか? ビジネステーブルを直接監視することは単純なケースでは機能しますが、トランザクショナルアウトボックスパターンは専用の、疎結合なイベントストリームを提供します。これにより、明示的なイベント契約を定義し、追加のメタデータでイベントをエンリッチし、公開するイベントペイロードが単なる生のデータベース行の変更ではなく、意図したとおりのものであることを保証できます。また、一貫したイベント構造を提供することで、コンシューマロジックを簡素化します。
-
DebeziumはPostgreSQLテーブルのスキーマ変更をどのように処理しますか? スキーマが変更された場合(例:列の追加、型の変更)、Debeziumはこれを検出します。AvroコンバータとSchema Registryで設定されている場合、そのトピックのAvroスキーマの新しいバージョンを登録します。その後、コンシューマはSchema Registryを使用して各メッセージの正しいスキーマバージョンを取得し、Schema Registryの互換性設定に応じて後方互換性と前方互換性を可能にします。
-
Debeziumコネクタがダウンした場合、データは失われますか? いいえ、データは失われません。PostgreSQLの論理レプリケーションスロットは、Debeziumが正常に処理するまでWALセグメントが保持されることを保証します。Debeziumが再起動すると、オフセットストレージの最後にコミットされたLSN(Log Sequence Number)から再開し、中断したところから正確に処理を続行します。ただし、長時間のダウンタイムはPostgreSQLサーバー上のWALの肥大化につながる可能性があります。
-
outboxテーブルから古いイベントをクリーンアップするにはどうすればよいですか? Debeziumはoutboxテーブルから読み取るだけで、レコードを削除することはありません。outboxテーブルから古いイベントを定期的に削除するには、別のプロセス(例:スケジュールされたジョブ)が必要です。このジョブは、十分に古く、Debeziumによって処理されたことが確認されたイベント(例:Debeziumのオフセットまたはタイムスタンプを確認する)のみを削除する必要があります。Debeziumがまだストリーミングしていないイベントを削除しないように注意してください。一般的な戦略は、特定のしきい値(例:7日)よりも古いイベントを削除することです。 -
pgoutput以外の論理デコーディングプラグインでもDebeziumを使用できますか? はい、Debeziumはwal2jsonなどの他のプラグインもサポートしています。ただし、pgoutputはPostgreSQL 10で導入されたネイティブの論理デコーディングプラグインであり、その効率性とPostgreSQLのレプリケーションプロトコルとの密接な統合から、一般的に推奨されています。wal2jsonはより人間が読みやすいJSON出力を提供し、デバッグや特定のユースケースに役立つかもしれませんが、pgoutputはDebeziumを使用した本番CDCの標準です。
Free In-Browser Developer Tools
Clean AI CLI logs, build cron expressions, decode JWTs, and calculate chmod permissions offline.
Related Articles
PostgreSQLのVACUUMとインデックス肥大化:検知、軽減、そして自動チューニング
PostgreSQLのテーブルとインデックスの肥大化を診断・解消します。自動バキュームのチューニング方法、pg_repackによるゼロダウンタイムでの再構築、MVCCの可視性マップまでを解説します。
Read more
PostgreSQL17クエリ最適化:実行プラン、メモリチューニング、EXPLAINANALYZE
PostgreSQL17のクエリ最適化について、実行プラン、メモリチューニング、EXPLAINANALYZEを網羅的に解説し、本番環境レベルのアーキテクチャとコード例を紹介する包括的なガイドです。
Read more
分散システムにおける分散ロック: Redlock、PostgreSQLアドバイザリロック、etcdリース
分散システムにおける分散ロックについて、Redlock、PostgreSQLアドバイザリロック、etcdリースを網羅し、本番環境レベルのアーキテクチャとコード例で解説する包括的なガイドです。
Read more