Konektor Debezium PolarDBO kompatibel dengan PolarDB for PostgreSQL (Compatible with Oracle) dan menangkap perubahan tingkat baris di database PolarDB for PostgreSQL (Compatible with Oracle), menghasilkan catatan event perubahan data, serta mengalirkannya ke topik Kafka. Untuk informasi lebih lanjut mengenai fitur dan penggunaannya, lihat konektor Debezium PostgreSQL dari komunitas.
Karena PolarDB for PostgreSQL (Compatible with Oracle) dan PostgreSQL komunitas hanya berbeda dalam penanganan beberapa tipe data dan objek bawaan, topik ini menjelaskan cara membangun konektor Debezium yang mendukung PolarDB for PostgreSQL (Compatible with Oracle) dengan menyesuaikan konektor Debezium PostgreSQL dari komunitas menggunakan perubahan kode minimal.
Build konektor Debezium PolarDBO
Konektor Debezium PolarDBO diadaptasi dari konektor Debezium PostgreSQL versi komunitas. Tidak tersedia perjanjian tingkat layanan (SLA) untuk konektor Debezium PolarDBO, baik Anda membangunnya sendiri maupun menggunakan paket JAR yang disediakan dalam topik ini.
Prasyarat
-
Siapkan lingkungan Java
Semua versi Debezium memerlukan Java 11 atau yang lebih baru. Pastikan Anda telah mengonfigurasi Java 11 sebelum membangun dan menjalankan konektor.
-
Tentukan versi Debezium
Pilih versi Debezium yang kompatibel dengan versi Kafka, Kafka Connect, dan PolarDB for PostgreSQL (Compatible with Oracle) Anda. Untuk detail kompatibilitas, lihat Ikhtisar Rilis Debezium.
Catatan-
Untuk repositori kode Debezium, lihat Debezium.
-
Tabel berikut memetakan versi PolarDB for PostgreSQL (Compatible with Oracle) ke versi PostgreSQL komunitas yang kompatibel.
-
Oracle compatibility 2.0 sesuai dengan PostgreSQL komunitas versi 14.
-
Oracle compatibility 1.0 sesuai dengan PostgreSQL komunitas versi 11.
-
-
-
Identifikasi versi PgJDBC
Dalam file
pom.xmldari versi Debezium yang dipilih, cariversion.postgresql.driveruntuk mengidentifikasi versi PgJDBC.CatatanUntuk repositori kode PgJDBC, lihat PgJDBC.
Prosedur
Debezium edisi komunitas 2.6.2.Final mendukung Kafka Connect 2.x dan 3.x, serta PostgreSQL versi 10, 11, 12, 13, 14, 15, dan 16.
Langkah-langkah berikut menunjukkan cara membangun konektor berdasarkan Debezium 2.6.2.Final.
-
Clone repositori Debezium dan PgJDBC untuk versi yang diperlukan.
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 -
Salin file PgJDBC yang diperlukan ke direktori 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 -
Terapkan file patch untuk menyesuaikan konektor agar kompatibel dengan PolarDB for PostgreSQL (Compatible with Oracle).
git apply v2.6.2.Final-support-polardbo-v1.patchCatatan-
Unduh file patch kompatibilitas untuk konektor Debezium PolarDBO: v2.6.2.Final-support-polardbo-v1.patch.
-
Secara default, patch ini mengemas dependensi debezium-api, debezium-core, PgJDBC, dan protobuf-java ke dalam file JAR. Untuk mengecualikan dependensi tersebut, hapus dari file pom.xml.
-
-
Gunakan Maven untuk membangun konektor Debezium PolarDBO.
mvn clean package -pl :debezium-connector-postgres -DskipITs -Dquick # Setelah proses build selesai, Anda dapat menemukan paket JAR di direktori debezium-connector-postgres/target.Paket JAR untuk konektor Debezium PolarDBO yang dibangun dengan JDK 11 juga tersedia untuk diunduh: debezium-connector-postgres-polardbo-v1.0-2.6.2.Final.jar.
Penggunaan
Konektor Debezium PolarDBO membaca perubahan inkremental dari database PolarDB for PostgreSQL (Compatible with Oracle) menggunakan replikasi logis. Kondisi berikut harus dipenuhi sebelum Anda dapat menggunakan konektor:
-
Parameter kluster
wal_levelharus diatur kelogical. Pengaturan ini menambahkan informasi yang diperlukan untuk replikasi logis ke dalam write-ahead logging (WAL).CatatanAnda dapat mengatur parameter kluster
wal_leveldi Konsol. Untuk informasi lebih lanjut, lihat Konfigurasi parameter kluster. Mengubah parameter ini akan merestart kluster. Rencanakan operasi ini secara hati-hati sesuai kebutuhan bisnis Anda. -
Jalankan perintah
ALTER TABLE schema.table REPLICA IDENTITY FULL;untuk mengaturREPLICA IDENTITYsetiap tabel yang dilangganan menjadiFULL. Hal ini memastikan bahwa event insert dan update berisi nilai sebelumnya untuk semua kolom, sehingga menjamin konsistensi data.Catatan-
REPLICA IDENTITY adalah pengaturan tingkat tabel khusus PostgreSQL. Pengaturan ini menentukan apakah plug-in decoding logis menyertakan nilai sebelumnya dari kolom tabel terkait selama event INSERT dan UPDATE. Untuk informasi lebih lanjut mengenai nilai-nilai REPLICA IDENTITY, lihat REPLICA IDENTITY.
-
Mengatur
REPLICA IDENTITYtabel yang dilangganan keFULLmungkin memerlukan penguncian tabel, yang dapat memengaruhi layanan Anda. Rencanakan operasi ini sesuai kebutuhan bisnis Anda. Anda dapat menjalankan perintah berikut untuk memeriksa apakah pengaturan saat ini sudahFULL:SELECT relreplident = 'f' FROM pg_class WHERE relname = 'tablename';
-
-
Nilai parameter
max_wal_sendersdanmax_replication_slotsharus lebih besar daripada jumlah total slot replikasi yang sedang digunakan ditambah yang dibutuhkan oleh pekerjaan Kafka Anda. -
Gunakan akun istimewa atau akun standar yang memiliki izin LOGIN dan REPLICATION. Akun tersebut juga harus memiliki izin SELECT pada semua tabel yang dilangganan untuk snapshot awal.
-
Hanya terhubung ke
primary endpointkluster PolarDB. Replikasi logis tidak didukung padacluster endpoint. -
Atur parameter
connector.classkeio.debezium.connector.postgresql.PolarDBOConnector. -
Kami merekomendasikan Anda mengatur parameter
plugin.namekepgoutput. Jika tidak, parsing inkremental dapat menghasilkan teks rusak (garbled) untuk database yang tidak menggunakan encoding UTF-8. Untuk informasi lebih lanjut, lihat dokumentasi komunitas.
Contoh
Contoh berikut menjelaskan cara menggunakan konektor Debezium PolarDBO untuk menyinkronkan tabel t1 dan t2 dari database dbz_db di kluster PolarDB for PostgreSQL dengan Oracle compatibility 2.0 ke antrian pesan Kafka.
Prasyarat
-
Siapkan Kafka
-
Deploy sebuah
instanceKafka dan pastikan dapat diakses darihostKafka Connect. Anda juga dapat menggunakan ApsaraMQ for Kafka. Untuk informasi lebih lanjut, lihat Quick start. -
Buat topik bernama
pg_dbz_eventdi instance Kafka untuk menerima pesan.CatatanUntuk pengujian, Anda dapat membuat topik dengan satu
partitionagar lebih mudah dilihat. Untuk lingkungan produksi, buat topik multi-partisi.
-
-
Jalankan Kafka Connect secara lokal dalam
distributed modepada port 8083.-
Salin paket JAR Debezium PolarDBO connector ke direktori
plugin.pathKafka Connect.# Ganti ${plugin.path} dengan path aktual. mkdir ${plugin.path}/debezium-connector-polardbo cp debezium-connector-postgres-polardbo-v1.0-2.6.2.Final.jar ${plugin.path}/debezium-connector-polardbo
-
-
Siapkan PolarDB for PostgreSQL (Compatible with Oracle)
-
Di halaman pembelian kluster PolarDB, beli kluster PolarDB for PostgreSQL (Compatible with Oracle) versi 2.0.
-
Konfigurasikan kluster PolarDB agar memenuhi semua prasyarat yang tercantum di bagian Penggunaan.
-
Buat akun istimewa. Untuk informasi lebih lanjut, lihat Buat akun.
-
Dapatkan
primary endpointkluster. Untuk informasi lebih lanjut, lihat Lihat titik akhir koneksi. Jika kluster PolarDB dan instance Kafka Connect berada diavailability zoneyang sama, Anda dapat menggunakanprivate endpoint. Jika tidak, Anda harus memintapublic endpoint. Tambahkan alamat instance Kafka Connect kecluster whitelistPolarDB. Untuk informasi lebih lanjut, lihat Konfigurasi daftar putih kluster. -
Di Konsol, buat database bernama
dbz_db. Untuk informasi lebih lanjut, lihat Buat database. -
Jalankan pernyataan berikut untuk membuat tabel
t1dant2di databasedbz_dbdan memasukkan data.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());
-
Tes
-
Buat file konfigurasi bernama
config/postgresql-connector.json. Untuk deskripsi parameter, lihat dokumentasi komunitas.{ "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" } }CatatanSecara default, Debezium membuat satu topik per tabel. Konfigurasi ini mengarahkan semua event perubahan ke satu topik tunggal.
-
Tambahkan konektor.
curl -i -X POST -H "Accept:application/json" -H "Content-Type:application/json" 'http://localhost:8083/connectors' -d @config/postgresql-connector.jsonSetelah konektor ditambahkan, data lengkap tersedia di topik Kafka.
Di tab message query instance Kafka Anda, atur metode kueri ke query by offset, pilih partisi
0, dan atur start offset ke0. Lalu, klik Query. Hasilnya menunjukkan dua pesan CDC Debezium (pada offset 0 dan 1). Kunci berisi nilai__dbz__physicalTableIdentifiermasing-masingpolardb.public.t1danpolardb.public.t2. Hal ini mengonfirmasi bahwa snapshot data lengkap telah disinkronkan ke Kafka. -
Jalankan pernyataan DML berikut di database
dbz_dbkluster PolarDB: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());Data inkremental kini tersedia di topik Kafka.
Di halaman kueri pesan, atur metode kueri ke query by offset, pilih partisi target dan offset awal, lalu klik Query. Hasilnya menunjukkan bahwa Debezium telah menangkap operasi DML pada tabel t1 dan t2 sebagai pesan CDC. Kunci setiap pesan mencakup pengenal tabel (seperti
polardb.public.t1) dan nilai primary key. Nilai pesan berisi event perubahan dalam format Debezium, termasuk bidangbeforedanafter. Untuk operasi DELETE, nilai pesan berukuran 0 byte.