Konektor Flink CDC untuk PolarDB for PostgreSQL (Compatible with Oracle), yang disebut sebagai konektor PolarDB-O Flink CDC, membaca snapshot data lengkap lalu perubahan inkremental dari database PolarDB for PostgreSQL (Compatible with Oracle). Untuk fitur dan penggunaannya, lihat dokumentasi komunitas Postgres CDC.
Karena PolarDB for PostgreSQL (Compatible with Oracle) dan PostgreSQL komunitas hanya memiliki perbedaan kecil dalam tipe data dan penanganan objek bawaan, artikel ini menjelaskan cara menyesuaikan konektor Postgres CDC komunitas dengan sedikit perubahan kode untuk mengemas konektor PolarDB Flink CDC untuk PolarDB for PostgreSQL (Compatible with Oracle).
Tipe DATE di PolarDB for PostgreSQL (Compatible with Oracle) berupa 64-bit, sedangkan di PostgreSQL komunitas berupa 32-bit. Konektor PolarDB-O Flink CDC menangani perbedaan ini.
Build PolarDB-O Flink CDC
Konektor PolarDB-O Flink CDC diadaptasi dari Postgres CDC komunitas. Tidak tersedia perjanjian tingkat layanan (SLA) untuk konektor ini, baik Anda membangunnya sendiri maupun menggunakan paket JAR yang disediakan dalam topik ini.
Prasyarat
-
Tentukan versi Flink-CDC
Jika Anda menggunakan Realtime Compute for Apache Flink Alibaba Cloud, Anda harus menentukan versi Flink-CDC komunitas yang kompatibel dengan versi Ververica Runtime (VVR) Anda. Untuk informasi lebih lanjut, lihat Pemetaan Versi CDC dan VVR.
CatatanRepositori kode Flink-CDC tersedia di Flink-CDC.
-
Tentukan versi Debezium
Dalam file
pom.xmldari versi Flink-CDC yang sesuai, temukan propertidebezium.versionuntuk menentukan versi Debezium.CatatanRepositori kode Debezium tersedia di Debezium.
-
Tentukan versi PgJDBC
Dalam file
pom.xmldari versi Postgres-CDC yang sesuai, temukan dependensiorg.postgresqluntuk menentukan versi PgJDBC.Catatan-
Untuk versi sebelum release-3.0, jalur file-nya adalah
flink-connector-postgres-cdc/pom.xml. -
Untuk release-3.0 dan seterusnya, jalur file-nya adalah
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-postgres-cdc/pom.xml. -
Repositori kode PgJDBC tersedia di PgJDBC.
-
Prosedur
Build untuk release-3.5
Flink-CDC komunitas release-3.5 kompatibel dengan vvr-11.4-jdk11-flink-1.20 dari Realtime Compute for Apache Flink Alibaba Cloud.
Untuk membangun konektor PolarDB-O Flink CDC untuk versi ini, ikuti langkah-langkah berikut:
-
Clone repositori untuk versi Flink-CDC, Debezium, dan PgJDBC yang sesuai.
git clone -b release-3.5 --depth=1 https://github.com/apache/flink-cdc.git git clone -b REL42.7.3 --depth=1 https://github.com/pgjdbc/pgjdbc.git git clone -b v1.9.8.Final --depth=1 https://github.com/debezium/debezium.git -
Salin file yang diperlukan dari repositori Debezium dan PgJDBC ke direktori Flink-CDC.
mkdir -p flink-cdc/flink-cdc-connect/flink-cdc-source-connectors/flink-connector-postgres-cdc/src/main/java/org/postgresql/core/v3 mkdir -p flink-cdc/flink-cdc-connect/flink-cdc-source-connectors/flink-connector-postgres-cdc/src/main/java/org/postgresql/jdbc cp pgjdbc/pgjdbc/src/main/java/org/postgresql/core/v3/ConnectionFactoryImpl.java flink-cdc/flink-cdc-connect/flink-cdc-source-connectors/flink-connector-postgres-cdc/src/main/java/org/postgresql/core/v3 cp pgjdbc/pgjdbc/src/main/java/org/postgresql/core/v3/QueryExecutorImpl.java flink-cdc/flink-cdc-connect/flink-cdc-source-connectors/flink-connector-postgres-cdc/src/main/java/org/postgresql/core/v3 cp pgjdbc/pgjdbc/src/main/java/org/postgresql/jdbc/PgDatabaseMetaData.java flink-cdc/flink-cdc-connect/flink-cdc-source-connectors/flink-connector-postgres-cdc/src/main/java/org/postgresql/jdbc cp pgjdbc/pgjdbc/src/main/java/org/postgresql/core/Oid.java flink-cdc/flink-cdc-connect/flink-cdc-source-connectors/flink-connector-postgres-cdc/src/main/java/org/postgresql/core cp debezium/debezium-connector-postgres/src/main/java/io/debezium/connector/postgresql/TypeRegistry.java flink-cdc/flink-cdc-connect/flink-cdc-source-connectors/flink-connector-postgres-cdc/src/main/java/io/debezium/connector/postgresql -
Masuk ke direktori Flink-CDC dan terapkan perbaikan bug konversi timestamp dan perbaikan bug perilaku implisit (SELECT *). Perbaikan ini dijadwalkan untuk dimasukkan dalam rilis komunitas 3.6.
cd flink-cdc # Terapkan perbaikan bug untuk konversi timestamp, yang akan digabungkan ke rilis komunitas 3.6. git fetch origin 2f32836a783f80f295c9dce339c11afec2a32dc2 git cherry-pick 2f32836a783f80f295c9dce339c11afec2a32dc2 git fetch origin 0d86de24494a855c2d83f9b1052c2e888e182cb1 git cherry-pick 0d86de24494a855c2d83f9b1052c2e888e182cb1 -
Terapkan file patch untuk memastikan kompatibilitas dengan PolarDB for PostgreSQL (Compatible with Oracle).
git apply release-3.5_support_polardbo.patchCatatanAnda dapat mengunduh file patch yang digunakan dalam langkah ini di sini: release-3.5_support_polardbo.patch.
-
Gunakan Maven untuk membangun konektor PolarDB-O Flink CDC.
mvn clean install -DskipTests -Dcheckstyle.skip=true -Dspotless.check.skip # Setelah build selesai, paket JAR berada di direktori flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-postgres/target.
Paket JAR berikut untuk konektor PolarDB-O Flink CDC dibangun dengan JDK 11 dengan mengikuti langkah-langkah di atas: flink-cdc-pipeline-connector-polardbo-3.5-SNAPSHOT-20260212.jar.
Build untuk release-3.1
Community Flink-CDC rilis-3.1 kompatibel dengan vvr-8.0.x-flink-1.17 dari Alibaba Cloud Realtime Compute for Apache Flink.
Untuk membangun konektor PolarDB-O Flink CDC untuk versi ini, ikuti langkah-langkah berikut:
-
Clone repositori untuk versi Flink-CDC, Debezium, dan PgJDBC yang sesuai.
git clone -b release-3.1 --depth=1 https://github.com/apache/flink-cdc.git git clone -b REL42.5.1 --depth=1 https://github.com/pgjdbc/pgjdbc.git git clone -b v1.9.8.Final --depth=1 https://github.com/debezium/debezium.git -
Salin file yang diperlukan dari repositori Debezium dan PgJDBC ke direktori Flink-CDC.
mkdir -p flink-cdc/flink-cdc-connect/flink-cdc-source-connectors/flink-connector-postgres-cdc/src/main/java/org/postgresql/core/v3 mkdir -p flink-cdc/flink-cdc-connect/flink-cdc-source-connectors/flink-connector-postgres-cdc/src/main/java/org/postgresql/jdbc cp pgjdbc/pgjdbc/src/main/java/org/postgresql/core/v3/ConnectionFactoryImpl.java flink-cdc/flink-cdc-connect/flink-cdc-source-connectors/flink-connector-postgres-cdc/src/main/java/org/postgresql/core/v3 cp pgjdbc/pgjdbc/src/main/java/org/postgresql/core/v3/QueryExecutorImpl.java flink-cdc/flink-cdc-connect/flink-cdc-source-connectors/flink-connector-postgres-cdc/src/main/java/org/postgresql/core/v3 cp pgjdbc/pgjdbc/src/main/java/org/postgresql/jdbc/PgDatabaseMetaData.java flink-cdc/flink-cdc-connect/flink-cdc-source-connectors/flink-connector-postgres-cdc/src/main/java/org/postgresql/jdbc cp debezium/debezium-connector-postgres/src/main/java/io/debezium/connector/postgresql/TypeRegistry.java flink-cdc/flink-cdc-connect/flink-cdc-source-connectors/flink-connector-postgres-cdc/src/main/java/io/debezium/connector/postgresql -
Terapkan file patch untuk memastikan kompatibilitas dengan PolarDB for PostgreSQL (Compatible with Oracle).
git apply release-3.1_support_polardbo.patchCatatanAnda dapat mengunduh file patch yang digunakan dalam langkah ini di sini: release-3.1_support_polardbo.patch.
-
Gunakan Maven untuk membangun konektor PolarDB-O Flink CDC.
mvn clean install -DskipTests -Dcheckstyle.skip=true -Dspotless.check.skip -Drat.skip=true # Setelah build selesai, paket JAR berada di direktori flink-sql-connector-postgres-cdc/target.
Paket JAR berikut untuk konektor PolarDB-O Flink CDC dibangun dengan JDK 8 dengan mengikuti langkah-langkah di atas: flink-sql-connector-postgres-cdc-3.1-SNAPSHOT.jar.
Build untuk release-2.3
Flink-CDC komunitas release-2.3 kompatibel dengan vvr-4.0.15-flink-1.13 hingga vvr-6.0.2-flink-1.15 dari Realtime Compute for Apache Flink Alibaba Cloud.
Untuk membangun konektor PolarDB-O Flink CDC untuk versi ini, ikuti langkah-langkah berikut:
-
Clone repositori untuk versi Flink-CDC, Debezium, dan PgJDBC yang sesuai.
git clone -b release-2.3 --depth=1 https://github.com/apache/flink-cdc.git git clone -b REL42.2.26 --depth=1 https://github.com/pgjdbc/pgjdbc.git git clone -b v1.6.4.Final --depth=1 https://github.com/debezium/debezium.git -
Salin file yang diperlukan dari repositori Debezium dan PgJDBC ke direktori Flink-CDC.
mkdir -p flink-cdc/flink-connector-postgres-cdc/src/main/java/org/postgresql/core/v3 mkdir -p flink-cdc/flink-connector-postgres-cdc/src/main/java/org/postgresql/jdbc mkdir -p flink-cdc/flink-connector-postgres-cdc/src/main/java/io/debezium/connector/postgresql cp pgjdbc/pgjdbc/src/main/java/org/postgresql/core/v3/ConnectionFactoryImpl.java flink-cdc/flink-connector-postgres-cdc/src/main/java/org/postgresql/core/v3 cp pgjdbc/pgjdbc/src/main/java/org/postgresql/core/v3/QueryExecutorImpl.java flink-cdc/flink-connector-postgres-cdc/src/main/java/org/postgresql/core/v3 cp pgjdbc/pgjdbc/src/main/java/org/postgresql/jdbc/PgDatabaseMetaData.java flink-cdc/flink-connector-postgres-cdc/src/main/java/org/postgresql/jdbc cp debezium/debezium-connector-postgres/src/main/java/io/debezium/connector/postgresql/TypeRegistry.java flink-cdc/flink-connector-postgres-cdc/src/main/java/io/debezium/connector/postgresql -
Terapkan file patch untuk memastikan kompatibilitas dengan PolarDB for PostgreSQL (Compatible with Oracle).
git apply release-2.3_support_polardbo.patchCatatanAnda dapat mengunduh file patch yang digunakan dalam langkah ini di sini: release-2.3_support_polardbo.patch.
-
Gunakan Maven untuk membangun konektor PolarDB-O Flink CDC.
mvn clean install -DskipTests -Dcheckstyle.skip=true -Dspotless.check.skip -Drat.skip=true # Setelah build selesai, paket JAR berada di direktori flink-sql-connector-postgres-cdc/target.
Paket JAR berikut untuk konektor PolarDB-O Flink CDC dibangun menggunakan JDK 8 sesuai langkah-langkah di atas: flink-sql-connector-postgres-cdc-2.3-SNAPSHOT.jar.
Penggunaan
Konektor PolarDB-O Flink CDC membaca aliran data CDC dari database PolarDB for PostgreSQL (Compatible with Oracle) dengan menggunakan logical replication. Ini memerlukan hal-hal berikut:
-
Atur parameter
wal_levelkelogical. Pengaturan ini menambahkan informasi yang diperlukan untuk logical replication ke file write-ahead logging (WAL).CatatanAnda dapat mengatur parameter wal_level di Konsol. Untuk instruksi detail, lihat Setel parameter kluster. Memodifikasi parameter ini akan me-restart kluster. Jadwalkan operasi ini dengan hati-hati untuk meminimalkan dampak pada bisnis Anda.
-
Jalankan perintah
ALTER TABLE schema.table REPLICA IDENTITY FULL;untuk mengaturREPLICA IDENTITYtabel yang dilangganan keFULL. Ini memastikan bahwa event INSERT dan UPDATE berisi nilai sebelumnya dari semua kolom, yang diperlukan untuk konsistensi data.Catatan-
REPLICA IDENTITY adalah pengaturan tingkat tabel di PostgreSQL yang menentukan apakah plugin logical decoding menyertakan nilai lama kolom untuk event INSERT dan UPDATE. Untuk informasi lebih lanjut tentang nilai REPLICA IDENTITY, lihat REPLICA IDENTITY.
-
Mengatur
REPLICA IDENTITYtabel yang dilangganan keFULLmungkin memerlukan kunci tabel dan berdampak pada operasi bisnis Anda. Rencanakan dengan sesuai. Anda dapat menggunakan perintah berikut untuk memeriksa apakah konfigurasi saat ini adalahFULL:SELECT relreplident = 'f' FROM pg_class WHERE relname = 'tablename';
-
-
Pastikan bahwa nilai parameter max_wal_senders dan max_replication_slots melebihi jumlah slot replikasi yang sedang digunakan ditambah jumlah slot yang dibutuhkan oleh pekerjaan Flink.
-
Gunakan akun istimewa, atau akun yang memiliki izin LOGIN dan REPLICATION. Akun tersebut juga harus memiliki izin SELECT pada tabel yang dilangganan untuk menjalankan kueri snapshot awal.
-
Anda hanya dapat terhubung ke titik akhir utama kluster PolarDB. Titik akhir kluster tidak mendukung logical replication.
-
Rilis 3.5 dan seterusnya mendukung sinkronisasi tabel partisi dengan menentukan tabel induknya. Konfigurasi berikut diperlukan. Untuk detailnya, lihat dokumentasi komunitas Postgres CDC.
-
Atur opsi
scan.include-partitioned-tables.enabledketrue. -
Buat secara manual
PUBLICATIONdi database dengan opsipublish_via_partition_root=true. Lalu, gunakan parameterdebezium.publication.nameuntuk menentukantable-name. -
table-nameharus hanya menentukan tabel induk. Ekspresi reguler tidak boleh mencocokkan tabel anak, karena akan menyebabkan duplikasi data selama fase snapshot.
Selain itu, rilis 3.5 dan seterusnya mendukung konektor pipeline, yang memungkinkan pembacaan data snapshot dan inkremental serta menyediakan sinkronisasi database penuh end-to-end. Namun, konektor pipeline saat ini tidak mendukung perubahan skema. Untuk detailnya, lihat dokumentasi komunitas Postgres CDC Pipeline Connector.
-
PolarDB-O Flink CDC vs. Postgres CDC
Konektor PolarDB-O Flink CDC dibangun berdasarkan konektor Postgres CDC. Untuk sintaksis dan parameter, lihat dokumentasi Postgres CDC. Namun, terdapat perbedaan utama berikut:
-
Dalam klausa WITH, parameter 'connector' harus diatur ke nilai tetap
polardbo-cdc. -
PolarDB Flink CDC kompatibel dengan semua versi PolarDB for PostgreSQL, PolarDB for PostgreSQL (Compatible with Oracle) 1.0, dan PolarDB for PostgreSQL (Compatible with Oracle) 2.0.
CatatanJika Anda menggunakan PolarDB for PostgreSQL, kami merekomendasikan untuk langsung menggunakan konektor Postgres CDC komunitas.
-
Untuk kolom tipe
DATEdi PolarDB for PostgreSQL (Compatible with Oracle) 1.0 dan PolarDB for PostgreSQL (Compatible with Oracle) 2.0, tipe yang sesuai di tabel sumber dan sink Flink SQL harus ditentukan sebagaiTIMESTAMP. -
Kami merekomendasikan mengatur parameter
decoding.plugin.namekepgoutput. Jika tidak, database dengan encoding non-UTF-8 mungkin menghasilkan karakter acak selama parsing inkremental. Untuk informasi lebih lanjut, lihat dokumentasi komunitas.
Pemetaan tipe data
Pemetaan tipe data antara PolarDB for PostgreSQL dan Flink identik dengan PostgreSQL komunitas, kecuali untuk tipe DATE. Pemetaan lengkapnya adalah sebagai berikut:
|
Tipe sumber |
Tipe Flink |
|
SMALLINT |
SMALLINT |
|
INT2 |
|
|
SMALLSERIAL |
|
|
SERIAL2 |
|
|
INTEGER |
INT |
|
SERIAL |
|
|
BIGINT |
BIGINT |
|
BIGSERIAL |
|
|
REAL |
FLOAT |
|
FLOAT4 |
|
|
FLOAT8 |
DOUBLE |
|
DOUBLE PRECISION |
|
|
NUMERIC(p, s) |
DECIMAL(p, s) |
|
DECIMAL(p, s) |
|
|
BOOLEAN |
BOOLEAN |
|
DATE |
|
|
TIME [(p)] [WITHOUT TIMEZONE] |
TIME [(p)] [WITHOUT TIMEZONE] |
|
TIMESTAMP [(p)] [WITHOUT TIMEZONE] |
TIMESTAMP [(p)] [WITHOUT TIMEZONE] |
|
CHAR(n) |
STRING |
|
CHARACTER(n) |
|
|
VARCHAR(n) |
|
|
CHARACTER VARYING(n) |
|
|
TEXT |
|
|
BYTEA |
BYTES |
Contoh
Konektor sumber
Contoh ini menunjukkan cara menggunakan konektor PolarDB-O Flink CDC untuk menyinkronkan tabel shipments dari database flink_source ke tabel shipments_sink di database flink_sink pada kluster PolarDB for PostgreSQL (Compatible with Oracle) 2.0.
Contoh ini memberikan demonstrasi dasar menjalankan konektor PolarDB-O Flink CDC di PolarDB for PostgreSQL (Compatible with Oracle). Untuk penggunaan produksi, konfigurasikan parameter konektor berdasarkan kebutuhan bisnis Anda dengan merujuk ke dokumentasi komunitas Postgres CDC.
-
Prasyarat
-
Siapkan PolarDB for PostgreSQL (Compatible with Oracle)
-
Di halaman pembelian kluster PolarDB, beli kluster PolarDB for PostgreSQL (Compatible with Oracle) 2.0.
-
Lihat titik akhir utama kluster. Jika kluster PolarDB dan ruang kerja Realtime Compute for Apache Flink Anda berada di virtual private cloud (VPC) yang sama, Anda dapat menggunakan titik akhir pribadi. Jika tidak, Anda harus mengajukan dan menggunakan titik akhir publik.
-
Konfigurasikan daftar putih alamat IP untuk kluster. Tambahkan alamat IP instans Flink ke daftar putih kluster PolarDB.
-
Di Konsol, buat database sumber flink_source dan database tujuan flink_sink. Untuk detailnya, lihat Buat database.
-
Jalankan pernyataan berikut untuk membuat tabel shipments di database flink_source dan memasukkan data.
CREATE TABLE public.shipments ( shipment_id INT, order_id INT, origin TEXT, destination TEXT, is_arrived BOOLEAN, order_time DATE, PRIMARY KEY (shipment_id) ); ALTER TABLE public.shipments REPLICA IDENTITY FULL; INSERT INTO public.shipments SELECT 1, 1, 'test1', 'test1', false, now(); -
Jalankan pernyataan berikut untuk membuat tabel shipments_sink di database flink_sink.
CREATE TABLE public.shipments_sink ( shipment_id INT, order_id INT, origin TEXT, destination TEXT, is_arrived BOOLEAN, order_time TIMESTAMP, PRIMARY KEY (shipment_id) );
-
-
Siapkan Realtime Compute for Apache Flink
-
Login ke Konsol Realtime Compute dan beli instans Realtime Compute for Apache Flink. Untuk informasi lebih lanjut, lihat Aktifkan Realtime Compute for Apache Flink.
CatatanKami merekomendasikan membuat ruang kerja Realtime Compute for Apache Flink di Region dan VPC yang sama dengan kluster PolarDB. Hal ini memungkinkan Anda menggunakan titik akhir utama pribadi kluster PolarDB untuk koneksi.
-
Buat konektor kustom dan unggah paket PolarDB-O Flink CDC yang telah Anda bangun. Pilih debezium-json untuk Formats. Untuk detailnya, lihat Buat konektor kustom.
-
-
-
Buat pekerjaan Flink
-
Login ke Konsol Realtime Compute for Apache Flink dan buat draft SQL. Untuk detailnya, lihat Kembangkan draft SQL. Gunakan kode Flink SQL berikut dan ganti placeholder untuk titik akhir utama, port, username, dan password kluster PolarDB Anda.
CatatanTipe DATE di PolarDB for PostgreSQL (Compatible with Oracle) berupa 64-bit, sedangkan tipe DATE di Flink SQL dan kebanyakan database berupa 32-bit. Oleh karena itu, Anda harus memetakan kolom
DATEdi tabel sumber ke tipeTIMESTAMPdi tabel sumber dan sink Flink SQL. Jika tidak, pekerjaan akan gagal dengan error ketidakcocokan tipe, seperti"java.time.DateTimeException: Invalid value for EpochDay (valid values -365243219162 - 365241780471): 1720891573000".CREATE TEMPORARY TABLE shipments ( shipment_id INT, order_id INT, origin STRING, destination STRING, is_arrived BOOLEAN, order_time TIMESTAMP, PRIMARY KEY (shipment_id) NOT ENFORCED ) WITH ( 'connector' = 'polardbo-cdc', 'hostname' = '<yourHostname>', 'port' = '<yourPort>', 'username' = '<yourUserName>', 'password' = '<yourPassWord>', 'database-name' = 'flink_source', 'schema-name' = 'public', 'table-name' = 'shipments', 'decoding.plugin.name' = 'pgoutput', 'slot.name' = 'flink' ); CREATE TEMPORARY TABLE shipments_sink ( shipment_id INT, order_id INT, origin STRING, destination STRING, is_arrived BOOLEAN, order_time TIMESTAMP, PRIMARY KEY (shipment_id) NOT ENFORCED ) WITH ( 'connector' = 'jdbc', 'url' = 'jdbc:postgresql://<yourHostname>:<yourPort>/flink_sink', 'table-name' = 'shipments_sink', 'username' = '<yourUserName>', 'password' = '<yourPassWord>' ); INSERT INTO shipments_sink SELECT * FROM shipments; -
Deploy dan mulai pekerjaan.
Di bilah alat atas editor pekerjaan, klik Deploy.
Di panel navigasi sebelah kiri, pilih Operations Center > Deployments. Di daftar pekerjaan, temukan pekerjaan target dan klik Start di kolom Actions.
-
Uji dan verifikasi hasilnya.
-
Setelah pekerjaan dideploy dan dalam status berjalan, data dari tabel shipments disinkronkan ke tabel shipments_sink di database flink_sink.
SELECT * FROM public.shipments_sink;Hasil berikut dikembalikan:
shipment_id | order_id | origin | destination | is_arrived | order_time -------------+----------+--------+-------------+------------+--------------------- 1 | 1 | test1 | test1 | f | 2024-09-18 05:45:08 (1 row) -
Jalankan pernyataan DML pada tabel shipments di database flink_source. Perubahan tersebut disinkronkan secara real time.
INSERT INTO public.shipments SELECT 2, 2, 'test2', 'test2', false, now(); UPDATE public.shipments SET is_arrived = true WHERE shipment_id = 1; DELETE FROM public.shipments WHERE shipment_id = 2; INSERT INTO public.shipments SELECT 3, 3, 'test3', 'test3', false, now(); UPDATE public.shipments SET is_arrived = true WHERE shipment_id = 3;Data di tabel shipments disinkronkan ke tabel shipments_sink di database flink_sink.
SELECT * FROM public.shipments_sink;Hasil berikut dikembalikan:
shipment_id | order_id | origin | destination | is_arrived | order_time -------------+----------+--------+-------------+------------+--------------------- 1 | 1 | test1 | test1 | t | 2024-09-18 05:45:08 3 | 3 | test3 | test3 | t | 2024-09-18 07:33:23 (2 rows)
-
-
Konektor pipeline
Contoh ini menunjukkan cara menggunakan konektor pipeline PolarDB-O Flink CDC untuk menyinkronkan tabel shipments1 dan shipments2 dari kluster PolarDB for PostgreSQL (Compatible with Oracle) 2.0. Untuk debugging, sink menggunakan konektor Print. Di lingkungan produksi, pilih konektor sink yang sesuai berdasarkan kebutuhan bisnis Anda.
Contoh ini memberikan demonstrasi dasar menjalankan konektor PolarDB-O Flink CDC di PolarDB for PostgreSQL (Compatible with Oracle). Untuk penggunaan produksi, konfigurasikan parameter berdasarkan kebutuhan bisnis Anda dengan merujuk ke dokumentasi komunitas Postgres CDC Pipeline Connector.
-
Prasyarat
-
Siapkan PolarDB for PostgreSQL (Compatible with Oracle)
-
Di halaman pembelian kluster PolarDB, beli kluster PolarDB for PostgreSQL (Compatible with Oracle) 2.0.
-
Lihat titik akhir utama kluster. Jika kluster PolarDB dan ruang kerja Realtime Compute for Apache Flink Anda berada di VPC yang sama, Anda dapat menggunakan titik akhir pribadi. Jika tidak, Anda harus mengajukan dan menggunakan titik akhir publik.
-
Konfigurasikan daftar putih alamat IP untuk kluster. Tambahkan alamat IP instans Flink ke daftar putih kluster PolarDB.
-
Di Konsol, buat database sumber flink_source. Untuk detailnya, lihat Buat database.
-
Jalankan pernyataan berikut untuk membuat tabel shipments1 dan shipments2 di database flink_source dan memasukkan data.
CREATE TABLE public.shipments1 ( shipment_id INT, order_id INT, origin TEXT, destination TEXT, is_arrived BOOLEAN, order_time DATE, PRIMARY KEY (shipment_id) ); ALTER TABLE public.shipments1 REPLICA IDENTITY FULL; INSERT INTO public.shipments1 SELECT 1, 1, 'test1', 'test1', false, now(); CREATE TABLE public.shipments2 ( shipment_id INT, order_id INT, origin TEXT, destination TEXT, is_arrived BOOLEAN, order_time DATE, PRIMARY KEY (shipment_id) ); ALTER TABLE public.shipments2 REPLICA IDENTITY FULL; INSERT INTO public.shipments2 SELECT 1, 1, 'test1', 'test1', false, now();
-
-
Siapkan Realtime Compute for Apache Flink
Login ke Konsol Realtime Compute dan beli instans Realtime Compute for Apache Flink. Untuk informasi lebih lanjut, lihat Aktifkan Realtime Compute for Apache Flink.
CatatanKami merekomendasikan membuat ruang kerja Realtime Compute for Apache Flink di Region dan VPC yang sama dengan kluster PolarDB. Hal ini memungkinkan Anda menggunakan titik akhir utama pribadi kluster PolarDB untuk koneksi.
-
-
Buat pekerjaan Flink
-
Login ke Konsol Realtime Compute for Apache Flink dan buat draft ingesti data. Untuk detailnya, lihat Ingesti data Flink CDC. Gunakan konfigurasi ingesti data berikut dan ganti placeholder untuk titik akhir utama, port, username, dan password kluster PolarDB Anda.
source: type: polardbo name: PolarDB Oracle Source hostname: '<yourHostname>' port: '<yourPort>' username: '<yourUserName>' password: '<yourPassWord>' tables: flink_source.public.shipments[12] decoding.plugin.name: pgoutput slot.name: pgtest sink: type: values name: values Sink print.enabled: true -
Di bagian More di sebelah kiri, tambahkan konektor pipeline yang telah Anda bangun. Di panel More Configurations di sebelah kanan, pastikan versi engine adalah
vvr-11.5-jdk11-flink-1.20dan tambahkan file dependensi yang diperlukan di area Additional Dependency Files. -
Deploy dan mulai pekerjaan.
-
Klik Deploy di pojok kanan atas.
source: type: polardbo name: PolarDB Oracle Source hostname: xxx port: xxx username: xxx password: xxx tables: flink_source.public.shipments[12] decoding.plugin.name: pgoutput slot.name: pgtest sink: type: values name: values Sink print.enabled: true -
Buka halaman Deployments dan klik Enable.
-
-
Uji dan verifikasi hasilnya.
-
Setelah pekerjaan deployment berhasil dijalankan, statusnya menjadi Running. Anda dapat menemukan CreateTableEvent dan DataChangeEvent dari fase data lengkap di log Job Log > Running Task Managers > Stdout. Di halaman detail deployment, klik tab Job Log, pilih Running Task Managers, lalu klik tab Stdout untuk melihat output log. Pastikan output berisi event pembuatan tabel dan event perubahan data berikut:
CreateTableEvent{tableId=public.shipments2, schema=columns={`shipment_id` INT NOT NULL,`order_id` INT,`origin` STRING,`destination` STRING,`is_arrived` BOOLEAN,`order_time` TIMESTAMP(6)}, primaryKeys=shipment_id, options=()} CreateTableEvent{tableId=public.shipments1, schema=columns={`shipment_id` INT NOT NULL,`order_id` INT,`origin` STRING,`destination` STRING,`is_arrived` BOOLEAN,`order_time` TIMESTAMP(6)}, primaryKeys=shipment_id, options=()} DataChangeEvent{tableId=public.shipments2, before=[], after=[1, 1, test1, test1, false, 2026-01-07T16:30:44], op=INSERT, meta=()} DataChangeEvent{tableId=public.shipments1, before=[], after=[1, 1, test1, test1, false, 2026-01-07T16:30:44], op=INSERT, meta=()} -
Jalankan pernyataan DML pada tabel shipments1 dan shipments2 di database flink_source. Perubahan tersebut disinkronkan secara real time.
INSERT INTO public.shipments1 SELECT 2, 2, 'test2', 'test2', false, now(); UPDATE public.shipments1 SET is_arrived = true WHERE shipment_id = 1; DELETE FROM public.shipments1 WHERE shipment_id = 2; INSERT INTO public.shipments1 SELECT 3, 3, 'test3', 'test3', false, now(); UPDATE public.shipments1 SET is_arrived = true WHERE shipment_id = 3; INSERT INTO public.shipments2 SELECT 2, 2, 'test2', 'test2', false, now(); UPDATE public.shipments2 SET is_arrived = true WHERE shipment_id = 1; DELETE FROM public.shipments2 WHERE shipment_id = 2; INSERT INTO public.shipments2 SELECT 3, 3, 'test3', 'test3', false, now(); UPDATE public.shipments2 SET is_arrived = true WHERE shipment_id = 3; -
Anda dapat menemukan DataChangeEvent untuk fase inkremental di log Job Logs > Running Task Managers > Stdout:
DataChangeEvent{tableId=public.shipments1, before=[], after=[2, 2, test2, test2, false, 2026-01-07T16:44:50], op=INSERT, meta=()} DataChangeEvent{tableId=public.shipments1, before=[1, 1, test1, test1, false, 2026-01-07T16:30:44], after=[1, 1, test1, test1, true, 2026-01-07T16:30:44], op=UPDATE, meta=()} DataChangeEvent{tableId=public.shipments1, before=[2, 2, test2, test2, false, 2026-01-07T16:44:50], after=[], op=DELETE, meta=()} DataChangeEvent{tableId=public.shipments1, before=[], after=[3, 3, test3, test3, false, 2026-01-07T16:44:50], op=INSERT, meta=()} DataChangeEvent{tableId=public.shipments1, before=[3, 3, test3, test3, false, 2026-01-07T16:44:50], after=[3, 3, test3, test3, true, 2026-01-07T16:44:50], op=UPDATE, meta=()} DataChangeEvent{tableId=public.shipments2, before=[], after=[2, 2, test2, test2, false, 2026-01-07T16:44:50], op=INSERT, meta=()} DataChangeEvent{tableId=public.shipments2, before=[1, 1, test1, test1, false, 2026-01-07T16:30:44], after=[1, 1, test1, test1, true, 2026-01-07T16:30:44], op=UPDATE, meta=()} DataChangeEvent{tableId=public.shipments2, before=[2, 2, test2, test2, false, 2026-01-07T16:44:50], after=[], op=DELETE, meta=()} DataChangeEvent{tableId=public.shipments2, before=[], after=[3, 3, test3, test3, false, 2026-01-07T16:44:50], op=INSERT, meta=()} DataChangeEvent{tableId=public.shipments2, before=[3, 3, test3, test3, false, 2026-01-07T16:44:50], after=[3, 3, test3, test3, true, 2026-01-07T16:44:50], op=UPDATE, meta=()}
-
-