Debezium PolarDBO コネクターは PolarDB for PostgreSQL (Compatible with Oracle) と互換性があり、PolarDB for PostgreSQL (Compatible with Oracle) データベースの行レベルの変更をキャプチャし、データ変更イベントレコードを生成して Kafka トピックにストリーミングします。その機能と使用方法の詳細については、コミュニティの「Debezium PostgreSQL コネクター」をご参照ください。
PolarDB for PostgreSQL (Compatible with Oracle) とコミュニティ版 PostgreSQL の違いは、少数のデータ型と組み込みオブジェクトの処理方法のみであるため、本トピックでは、コミュニティの Debezium PostgreSQL コネクターを最小限のコード変更で適応させ、PolarDB for PostgreSQL (Compatible with Oracle) をサポートする Debezium コネクターをビルドする方法について説明します。
Debezium PolarDBO コネクターのビルド
Debezium PolarDBO コネクターは、コミュニティの Debezium PostgreSQL コネクターを基にしています。ご自身でビルドした場合でも、本トピックで提供される JAR パッケージを使用した場合でも、Debezium PolarDBO コネクターにはサービスレベルアグリーメント (SLA) は提供されません。
前提条件
-
Java 環境のセットアップ
すべての Debezium バージョンには Java 11 以降が必要です。コネクターをビルドして実行する前に、Java 11 が設定されていることを確認してください。
-
Debezium バージョンの決定
お使いの Kafka、Kafka Connect、および PolarDB for PostgreSQL (Compatible with Oracle) のバージョンと互換性のある Debezium バージョンを選択してください。互換性の詳細については、「Debezium リリース概要」をご参照ください。
説明-
Debezium のコードリポジトリについては、「Debezium」をご参照ください。
-
次の表は、PolarDB for PostgreSQL (Compatible with Oracle) のバージョンと互換性のあるコミュニティ版 PostgreSQL のバージョンを示しています。
-
Oracle compatibility 2.0 は、コミュニティ版 PostgreSQL 14 に対応します。
-
Oracle compatibility 1.0 は、コミュニティ版 PostgreSQL 11 に対応します。
-
-
-
PgJDBC バージョンの特定
選択した Debezium バージョンの
pom.xmlファイルでversion.postgresql.driverを検索し、PgJDBC のバージョンを特定します。説明PgJDBC のコードリポジトリについては、「PgJDBC」をご参照ください。
手順
Debezium コミュニティ版 2.6.2.Final は、Kafka Connect 2.x および 3.x、ならびに PostgreSQL バージョン 10、11、12、13、14、15、16 をサポートしています。
以下の手順では、Debezium 2.6.2.Final に基づいてコネクターをビルドする方法を説明します。
-
必要なバージョンの Debezium と PgJDBC のリポジトリをクローンします。
git clone -b v2.6.2.Final --depth=1 https://github.com/debezium/debezium.git git clone -b REL42.6.1 --depth=1 https://github.com/pgjdbc/pgjdbc.git -
必要な PgJDBC ファイルを Debezium ディレクトリにコピーします。
mkdir -p debezium/debezium-connector-postgres/src/main/java/org/postgresql/core/v3 mkdir -p debezium/debezium-connector-postgres/src/main/java/org/postgresql/jdbc cp pgjdbc/pgjdbc/src/main/java/org/postgresql/core/v3/ConnectionFactoryImpl.java debezium/debezium-connector-postgres/src/main/java/org/postgresql/core/v3 cp pgjdbc/pgjdbc/src/main/java/org/postgresql/core/v3/QueryExecutorImpl.java debezium/debezium-connector-postgres/src/main/java/org/postgresql/core/v3 cp pgjdbc/pgjdbc/src/main/java/org/postgresql/jdbc/PgDatabaseMetaData.java debezium/debezium-connector-postgres/src/main/java/org/postgresql/jdbc -
パッチファイルを適用して、コネクターを PolarDB for PostgreSQL (Compatible with Oracle) に対応させます。
git apply v2.6.2.Final-support-polardbo-v1.patch説明-
Debezium PolarDBO コネクターの互換性パッチファイルをダウンロードします:v2.6.2.Final-support-polardbo-v1.patch。
-
デフォルトでは、このパッチは debezium-api、debezium-core、PgJDBC、および protobuf-java の依存関係を JAR にパッケージ化します。これらの依存関係を除外するには、pom.xml ファイルから削除してください。
-
-
Maven を使用して Debezium PolarDBO コネクターをビルドします。
mvn clean package -pl :debezium-connector-postgres -DskipITs -Dquick # ビルドが完了すると、JAR パッケージは debezium-connector-postgres/target ディレクトリに生成されます。JDK 11 でビルドされた Debezium PolarDBO コネクターの JAR パッケージもダウンロードできます:debezium-connector-postgres-polardbo-v1.0-2.6.2.Final.jar。
使用方法
Debezium PolarDBO コネクターは、論理レプリケーションを使用して PolarDB for PostgreSQL (Compatible with Oracle) データベースから増分変更を読み取ります。コネクターを使用する前に、以下の条件を満たす必要があります:
-
wal_levelクラスターパラメーターをlogicalに設定する必要があります。この設定により、論理レプリケーションに必要な情報が先行書き込みログ (WAL) に追加されます。説明wal_levelクラスターパラメーターはコンソールで設定できます。詳細については、「クラスターパラメーターの設定」をご参照ください。このパラメーターを変更するとクラスターが再起動します。ビジネスニーズに応じて、この操作を慎重に計画してください。 -
ALTER TABLE schema.table REPLICA IDENTITY FULL;コマンドを実行して、各サブスクライブされたテーブルのREPLICA IDENTITYをFULLに設定します。これにより、挿入および更新イベントにすべての列の以前の値が含まれるようになり、データの一貫性が保証されます。説明-
REPLICA IDENTITY は PostgreSQL 固有のテーブルレベルの設定です。これにより、INSERT および UPDATE イベント中に論理デコーディングプラグインが関連するテーブル列の以前の値を含めるかどうかが決まります。REPLICA IDENTITY の値の詳細については、「REPLICA IDENTITY」をご参照ください。
-
サブスクライブされたテーブルの
REPLICA IDENTITYをFULLに設定すると、テーブルロックが必要になる場合があり、サービスに影響を与える可能性があります。ビジネスニーズに応じて、この操作を計画してください。次のコマンドを実行して、現在の設定がFULLであるかどうかを確認できます:SELECT relreplident = 'f' FROM pg_class WHERE relname = 'tablename';
-
-
max_wal_sendersおよびmax_replication_slotsパラメーターの値は、現在使用中のレプリケーションスロットと Kafka ジョブで必要なレプリケーションスロットの合計よりも大きい必要があります。 -
特権アカウント、または LOGIN と REPLICATION の両方の権限を持つ標準アカウントを使用します。初期スナップショットのために、アカウントはすべてのサブスクライブされたテーブルに対する SELECT 権限も持っている必要があります。
-
PolarDB クラスターの
プライマリエンドポイントにのみ接続してください。論理レプリケーションはクラスターエンドポイントではサポートされていません。 -
connector.classパラメーターをio.debezium.connector.postgresql.PolarDBOConnectorに設定します。 -
plugin.nameパラメーターをpgoutputに設定することを推奨します。そうしないと、UTF-8 エンコーディングを使用しないデータベースで増分解析を行うと、テキストが文字化けする可能性があります。詳細については、「コミュニティドキュメント」をご参照ください。
例
次の例では、Debezium PolarDBO コネクターを使用して、PolarDB for PostgreSQL Oracle 互換性 2.0 クラスター内の dbz_db データベースから t1 テーブルと t2 テーブルを Kafka メッセージキューに同期する方法について説明します。
前提条件
-
Kafka のセットアップ
-
Kafka
インスタンスをデプロイし、Kafka Connectホストからアクセスできることを確認します。ApsaraMQ for Kafka を使用することもできます。詳細については、「クイックスタート」をご参照ください。 -
メッセージを受信するために、Kafka インスタンスに
pg_dbz_eventという名前のトピックを作成します。説明テスト目的では、表示を容易にするために単一
パーティションのトピックを作成できます。本番環境では、複数パーティションのトピックを作成してください。
-
-
ポート 8083 で、Kafka Connect をローカルで
分散モードで起動します。-
Debezium PolarDBO コネクターの JAR パッケージを Kafka Connect の
plugin.pathディレクトリにコピーします。# ${plugin.path} を実際のパスに置き換えます。 mkdir ${plugin.path}/debezium-connector-polardbo cp debezium-connector-postgres-polardbo-v1.0-2.6.2.Final.jar ${plugin.path}/debezium-connector-polardbo
-
-
PolarDB for PostgreSQL (Compatible with Oracle) のセットアップ
-
PolarDB クラスター購入ページで、PolarDB for PostgreSQL (Compatible with Oracle) 2.0 クラスターを購入します。
-
使用方法セクションに記載されているすべての前提条件を満たすように PolarDB クラスターを設定します。
-
特権アカウントを作成します。詳細については、「アカウントの作成」をご参照ください。
-
クラスターの
プライマリエンドポイントを取得します。詳細については、「接続エンドポイントの表示」をご参照ください。PolarDB クラスターと Kafka Connectインスタンスのアドレスが同じアベイラビリティーゾーンにある場合は、プライベートエンドポイントを使用できます。それ以外の場合は、パブリックエンドポイントをリクエストする必要があります。Kafka Connectインスタンスのアドレスを PolarDBクラスターホワイトリストに追加します。詳細については、「クラスターホワイトリストの設定」をご参照ください。 -
コンソールで、
dbz_dbという名前のデータベースを作成します。詳細については、「データベースの作成」をご参照ください。 -
次のステートメントを実行して、
dbz_dbデータベースにt1およびt2テーブルを作成し、データを挿入します。CREATE TABLE public.t1 (a int PRIMARY KEY, b text, c TIMESTAMP); ALTER TABLE public.t1 REPLICA IDENTITY FULL; INSERT INTO public.t1(a, b, c) VALUES(1, 'a', now()); CREATE TABLE public.t2 (a int PRIMARY KEY, b text, c DATE); ALTER TABLE public.t2 REPLICA IDENTITY FULL; INSERT INTO public.t2(a, b, c) VALUES(1, 'a', now());
-
テスト
-
config/postgresql-connector.jsonという名前の設定ファイルを作成します。パラメーターの説明については、「コミュニティドキュメント」をご参照ください。{ "name": "dbz-polardb", "config": { "connector.class": "io.debezium.connector.postgresql.PolarDBOConnector", "database.hostname": "<yourHostname>", "database.port": "<yourPort>", "database.user": "<yourUserName>", "database.password": "<yourPassWord>", "database.dbname" : "dbz_db", "plugin.name": "pgoutput", "slot.name": "dbz_polardb", "table.include.list": "public.t1,public.t2", "topic.prefix": "polardb", "transforms": "Combine", "transforms.Combine.type": "io.debezium.transforms.ByLogicalTableRouter", "transforms.Combine.topic.regex": "(.*)", "transforms.Combine.topic.replacement": "pg_dbz_event" } }説明デフォルトでは、Debezium はテーブルごとに 1 つのトピックを作成します。この設定では、すべての変更イベントが単一のトピックにルーティングされます。
-
コネクターを追加します。
curl -i -X POST -H "Accept:application/json" -H "Content-Type:application/json" 'http://localhost:8083/connectors' -d @config/postgresql-connector.jsonコネクターが追加されると、完全なデータが Kafka トピックで利用できるようになります。
Kafka インスタンスの [メッセージクエリ] タブで、クエリ方法を [オフセットによるクエリ] に設定し、パーティション
0を選択して、[開始オフセット] を0に設定します。次に、[クエリ] をクリックします。結果には、2 つの Debezium CDC メッセージ (オフセット 0 と 1) が表示されます。キーには、それぞれ__dbz__physicalTableIdentifierの値がpolardb.public.t1とpolardb.public.t2として含まれています。これにより、完全なデータスナップショットが Kafka に同期されたことが確認できます。 -
PolarDB クラスターの
dbz_dbデータベースで、次の DML ステートメントを実行します:INSERT INTO public.t1(a, b, c) VALUES(2, 'b', now()); UPDATE public.t1 SET b = 'c' WHERE a = 1; DELETE FROM public.t1 WHERE a = 2; INSERT INTO public.t1(a, b, c) VALUES(4, 'd', now()); INSERT INTO public.t2(a, b, c) VALUES(2, 'b', now()); UPDATE public.t2 SET b = 'c' WHERE a = 1; DELETE FROM public.t2 WHERE a = 2; INSERT INTO public.t2(a, b, c) VALUES(4, 'd', now());増分データが Kafka トピックで利用できるようになります。
メッセージクエリページで、 [クエリ方法] を [オフセットによるクエリ] に設定し、対象のパーティションと開始オフセットを選択して、 [クエリ] をクリックします。結果には、Debezium が t1 および t2 テーブルに対する DML 操作を CDC メッセージとしてキャプチャしたことが示されます。各メッセージのキーには、テーブル識別子 (例:
polardb.public.t1) と主キーの値が含まれています。メッセージの値には、beforeおよびafterフィールドを含む、Debezium 形式の変更イベントが含まれています。DELETE 操作の場合、afterフィールドはnullになり、beforeフィールドに削除された行のデータが格納されます。