All Products
Search
Document Center

Realtime Compute for Apache Flink:Akses AnalyticDB for PostgreSQL

Last Updated:Jul 25, 2026

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:

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

  1. Masuk ke Konsol AnalyticDB for PostgreSQL.

  2. Tambahkan Blok CIDR dari ruang kerja Flink yang sepenuhnya dikelola ke daftar putih instans AnalyticDB for PostgreSQL.

    1. Temukan Blok CIDR dari vSwitch yang digunakan oleh ruang kerja Flink yang sepenuhnya dikelola. Lihat How do I configure a whitelist?

    2. 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.
  3. 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.

  4. Buat tabel dimensi bernama adbpg_dim_table dan 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.

  5. Buat tabel hasil bernama adbpg_sink_table agar Flink dapat menulis output ke sana.

    CREATE TABLE adbpg_sink_table(
      id int,
      username text,
      score int
    );

Langkah 2: Buat draf aliran

  1. Masuk ke Konsol Realtime Compute for Apache Flink, temukan ruang kerja Anda, lalu klik Console pada kolom Actions.

  2. Di panel navigasi sebelah kiri, buka Development > ETL. Di pojok kiri atas halaman Editor SQL, klik +, lalu pilih New Blank Stream Draft.

  3. Pada kotak dialog New Draft, konfigurasikan parameter berikut.

    Parameter Deskripsi Contoh
    Name Nama draf. Harus unik dalam proyek. adbpg-test
    Location Folder tempat draf disimpan. Klik ikon di samping folder yang sudah ada untuk membuat subfolder. Draft
    Engine Version Versi mesin Flink. Lihat Engine versions untuk detail versi dan siklus hidup. vvr-8.0.1-flink-1.17
  4. Klik Create.

Langkah 3: Tulis dan terapkan draf

  1. 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.

  2. Perbarui parameter koneksi untuk tabel dimensi dan tabel hasil. Ganti nilai placeholder pada klausa WITH dengan 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
    url Ya URL JDBC dalam format jdbc:postgresql://<Internal endpoint>:<Port>/<Database name>. Temukan ini di halaman Database Connection instans di Konsol AnalyticDB for PostgreSQL.
    tablename Ya Nama tabel di database AnalyticDB for PostgreSQL.
    username Ya Username untuk mengakses database.
    password Ya Password untuk akun database.
    targetSchema Tidak public Nama skema. Tentukan hanya jika tabel Anda tidak berada di skema public.
    maxRetryTimes Tidak Jumlah maksimum percobaan ulang setelah kegagalan penulisan.
    cache Tidak Kebijakan cache untuk pencarian tabel dimensi. Atur ke lru untuk 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 — sesuaikan cacheSize dan pertimbangkan toleransi Anda terhadap data kedaluwarsa sebelum mengaktifkan fitur ini.
    cacheSize Tidak Jumlah maksimum entri yang di-cache. Nilai yang lebih besar mengurangi permintaan ke database tetapi mengonsumsi lebih banyak memori.
    conflictMode Tidak Tindakan yang diambil ketika penulisan bentrok dengan primary key atau indeks yang sudah ada. Atur ke ignore untuk melewatkan baris yang bentrok.
    retryWaitTime Tidak Waktu dalam milidetik untuk menunggu antar percobaan ulang penulisan.
  3. Di pojok kanan atas halaman Editor SQL, klik Validate untuk memeriksa sintaksis.

  4. Klik Deploy.

  5. Pada halaman O&M > Deployments, temukan penerapan Anda, lalu klik Start pada kolom Actions.

Langkah 4: Verifikasi hasil

  1. Masuk ke Konsol AnalyticDB for PostgreSQL.

  2. Klik Log On to Database. Untuk detail selengkapnya, lihat Connect to an instance from a client.

  3. 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.

Referensi