Tutorial ini menjelaskan cara menggunakan AnalyticDB for PostgreSQL sebagai tabel dimensi sekaligus tabel hasil dalam pekerjaan Flink SQL — pola umum untuk pipeline penguatan data real-time.
Di akhir tutorial ini, Anda akan memiliki pekerjaan Flink yang berjalan, yang membaca dari sumber Datagen, mencari informasi pengguna dari tabel dimensi AnalyticDB for PostgreSQL, lalu menulis catatan yang diperkaya ke tabel hasil AnalyticDB for PostgreSQL.
Batasan
-
Realtime Compute for Apache Flink tidak dapat membaca dari AnalyticDB for PostgreSQL dalam mode serverless.
-
Konektor AnalyticDB for PostgreSQL memerlukan Ververica Runtime (VVR) versi 6.0.0 atau lebih baru.
-
AnalyticDB for PostgreSQL V7.0 memerlukan VVR 8.0.1 atau lebih baru.
Untuk menggunakan konektor kustom sebagai gantinya, lihat Manage custom connectors.
Prasyarat
Sebelum memulai, pastikan Anda telah memiliki:
-
Ruang kerja Flink yang sepenuhnya dikelola. Lihat Activate fully managed Flink.
-
Instans AnalyticDB for PostgreSQL dan akun istimewa. Lihat Create an instance dan Create a privileged account.
-
Instans AnalyticDB for PostgreSQL dan ruang kerja Flink yang sepenuhnya dikelola berada dalam virtual private cloud (VPC) yang sama.
Jika keduanya berada di VPC yang berbeda, lihat How does fully managed Flink access a service across VPCs?
Langkah 1: Konfigurasi daftar putih dan siapkan data
-
Masuk ke Konsol AnalyticDB for PostgreSQL.
-
Tambahkan Blok CIDR dari ruang kerja Flink yang sepenuhnya dikelola ke daftar putih instans AnalyticDB for PostgreSQL.
-
Temukan Blok CIDR dari vSwitch yang digunakan oleh ruang kerja Flink yang sepenuhnya dikelola. Lihat How do I configure a whitelist?
-
Tambahkan Blok CIDR tersebut ke daftar putih instans AnalyticDB for PostgreSQL. Lihat Procedure.
Jika Anda mengakses instans melalui Internet, tambahkan Alamat IP publik sebagai gantinya.
-
-
Pada halaman detail instans, klik Log On to Database di pojok kanan atas, lalu masukkan username dan password Anda. Untuk detail selengkapnya, lihat Use client tools to connect to an instance.
-
Buat tabel dimensi bernama
adbpg_dim_tabledan masukkan 50 baris data sampel.-- Buat tabel dimensi CREATE TABLE adbpg_dim_table( id int, username text, PRIMARY KEY(id) ); -- Masukkan 50 baris: id berkisar dari 1 hingga 50, username adalah "username" diikuti nomor baris INSERT INTO adbpg_dim_table(id, username) SELECT i, 'username'||i::text FROM generate_series(1, 50) AS t(i);Jalankan
SELECT * FROM adbpg_dim_table ORDER BY id;untuk memverifikasi data yang dimasukkan. -
Buat tabel hasil bernama
adbpg_sink_tableagar Flink dapat menulis output ke sana.CREATE TABLE adbpg_sink_table( id int, username text, score int );
Langkah 2: Buat draf aliran
-
Masuk ke Konsol Realtime Compute for Apache Flink, temukan ruang kerja Anda, lalu klik Console pada kolom Actions.
-
Di panel navigasi sebelah kiri, buka Development > ETL. Di pojok kiri atas halaman Editor SQL, klik +, lalu pilih New Blank Stream Draft.
-
Pada kotak dialog New Draft, konfigurasikan parameter berikut.
Parameter Deskripsi Contoh Name Nama draf. Harus unik dalam proyek. adbpg-testLocation Folder tempat draf disimpan. Klik ikon di samping folder yang sudah ada untuk membuat subfolder. DraftEngine Version Versi mesin Flink. Lihat Engine versions untuk detail versi dan siklus hidup. vvr-8.0.1-flink-1.17 -
Klik Create.
Langkah 3: Tulis dan terapkan draf
-
Salin SQL berikut ke editor kode. SQL ini mendefinisikan tiga tabel dan join lookup yang memperkaya aliran Datagen dengan data pengguna dari AnalyticDB for PostgreSQL.
-- Tabel sumber: Datagen menghasilkan ID berurutan (1-50) dan skor acak (70-100). -- Tidak perlu perubahan pada klausa WITH untuk contoh ini. CREATE TEMPORARY TABLE datagen_source ( id INT, score INT ) WITH ( 'connector' = 'datagen', 'fields.id.kind' = 'sequence', 'fields.id.start' = '1', 'fields.id.end' = '50', 'fields.score.kind' = 'random', 'fields.score.min' = '70', 'fields.score.max' = '100' ); -- Tabel dimensi: didukung oleh AnalyticDB for PostgreSQL. -- Flink melakukan kueri ke tabel ini pada waktu pemrosesan untuk mencari username berdasarkan ID. -- Ganti nilai pada klausa WITH dengan detail koneksi aktual Anda. CREATE TEMPORARY TABLE dim_adbpg( id int, username varchar, PRIMARY KEY(id) NOT ENFORCED ) WITH ( 'connector' = 'adbpg', 'url' = 'jdbc:postgresql://gp-2ze****3tysk255b5-master.gpdb.rds.aliyuncs.com:5432/flinktest', 'tablename' = 'adbpg_dim_table', 'username' = 'flinktest', 'password' = '${secret_values.adb_password}', 'maxRetryTimes' = '2', 'cache' = 'lru', 'cacheSize' = '100' ); -- Tabel hasil: Flink menulis catatan yang diperkaya ke sini. -- Ganti nilai pada klausa WITH dengan detail koneksi aktual Anda. CREATE TEMPORARY TABLE sink_adbpg ( id int, username varchar, score int ) WITH ( 'connector' = 'adbpg', 'url' = 'jdbc:postgresql://gp-2ze****3tysk255b5-master.gpdb.rds.aliyuncs.com:5432/flinktest', 'tablename' = 'adbpg_sink_table', 'username' = 'flinktest', 'password' = '${secret_values.adb_password}', 'maxRetryTimes' = '2', 'conflictMode' = 'ignore', 'retryWaitTime' = '200' ); -- Lookup join: untuk setiap catatan dari datagen_source, Flink mencari baris yang sesuai -- di dim_adbpg pada saat catatan diproses (PROCTIME()). INSERT INTO sink_adbpg SELECT ts.id, ts.username, ds.score FROM datagen_source AS ds JOIN dim_adbpg FOR SYSTEM_TIME AS OF PROCTIME() AS ts ON ds.id = ts.id;Mengenai sintaksis lookup join:
FOR SYSTEM_TIME AS OF PROCTIME()memberi tahu Flink untuk mencari tabel dimensi pada saat setiap catatan sumber diproses. Artinya, setiap catatan diperkaya dengan data dimensi yang tersedia pada waktu pemrosesan, dan hasil yang sudah ditulis tidak diperbarui meskipun tabel dimensi berubah di kemudian hari. -
Perbarui parameter koneksi untuk tabel dimensi dan tabel hasil. Ganti nilai placeholder pada klausa
WITHdengan detail koneksi AnalyticDB for PostgreSQL Anda yang sebenarnya. Tabel sumber Datagen tidak memerlukan perubahan. Untuk referensi lengkap mengenai parameter dan pemetaan tipe data, lihat AnalyticDB for PostgreSQL connector.Parameter Wajib Bawaan Deskripsi urlYa — URL JDBC dalam format jdbc:postgresql://<Internal endpoint>:<Port>/<Database name>. Temukan ini di halaman Database Connection instans di Konsol AnalyticDB for PostgreSQL.tablenameYa — Nama tabel di database AnalyticDB for PostgreSQL. usernameYa — Username untuk mengakses database. passwordYa — Password untuk akun database. targetSchemaTidak publicNama skema. Tentukan hanya jika tabel Anda tidak berada di skema public.maxRetryTimesTidak — Jumlah maksimum percobaan ulang setelah kegagalan penulisan. cacheTidak — Kebijakan cache untuk pencarian tabel dimensi. Atur ke lruuntuk menyimpan entri yang baru saja diakses di memori. Caching LRU mengurangi traffic database dan meningkatkan throughput lookup, tetapi entri yang di-cache mungkin kedaluwarsa. Ini merupakan pertukaran antara throughput dan kesegaran data — sesuaikancacheSizedan pertimbangkan toleransi Anda terhadap data kedaluwarsa sebelum mengaktifkan fitur ini.cacheSizeTidak — Jumlah maksimum entri yang di-cache. Nilai yang lebih besar mengurangi permintaan ke database tetapi mengonsumsi lebih banyak memori. conflictModeTidak — Tindakan yang diambil ketika penulisan bentrok dengan primary key atau indeks yang sudah ada. Atur ke ignoreuntuk melewatkan baris yang bentrok.retryWaitTimeTidak — Waktu dalam milidetik untuk menunggu antar percobaan ulang penulisan. -
Di pojok kanan atas halaman Editor SQL, klik Validate untuk memeriksa sintaksis.
-
Klik Deploy.
-
Pada halaman O&M > Deployments, temukan penerapan Anda, lalu klik Start pada kolom Actions.
Langkah 4: Verifikasi hasil
-
Masuk ke Konsol AnalyticDB for PostgreSQL.
-
Klik Log On to Database. Untuk detail selengkapnya, lihat Connect to an instance from a client.
-
Jalankan kueri berikut untuk melihat catatan yang ditulis Flink ke tabel hasil.
SELECT * FROM adbpg_sink_table ORDER BY id;Hasilnya harus berisi 50 baris, masing-masing dengan ID pengguna, username yang sesuai dari tabel dimensi, dan skor acak antara 70 hingga 100.
Kueri mengembalikan 7 catatan dengan tiga kolom: id, username, dan score. Baris 1 hingga 7 masing-masing sesuai dengan username1 hingga username7, dengan skor 94, 79, 70, 93, 71, 82, dan 87.