All Products
Search
Document Center

Data Lake Formation:Akses katalog DLF dari Flink CDC menggunakan Iceberg REST

Last Updated:May 28, 2026

Gunakan Flink CDC untuk menyinkronkan data ke katalog Data Lake Formation (DLF) melalui Iceberg REST dengan Realtime Compute for Apache Flink.

Prasyarat

Sebelum memulai, pastikan Anda telah memiliki:

Batasan

Konektivitas Iceberg REST ke DLF memerlukan mesin Realtime Compute for Apache Flink versi VVR 11.6.0 atau lebih baru.

Daftarkan katalog DLF di Flink

Pendaftaran katalog di Flink membuat pemetaan ke katalog DLF Anda. Membuat atau menghapus katalog di Flink tidak memengaruhi data aktual di DLF. Semua tabel yang dibuat di katalog DLF melalui Iceberg REST adalah tabel Iceberg.
  1. Masuk ke Konsol Manajemen Realtime Compute for Apache Flink.

  2. Pada kolom Actions ruang kerja Anda, klik Console.

  3. Di panel navigasi kiri, klik Development > Scripts.

  4. Buat skrip baru, lalu tempel pernyataan SQL berikut ke editor SQL.

    CREATE CATALOG `catalog_name`
     WITH (
        'type' = 'iceberg',
        'catalog-type' = 'rest',
        'uri' = 'http://{region-id}-vpc.dlf.aliyuncs.com/iceberg',
        'warehouse' = 'iceberg_test',
        'rest.signing-region' = '{region-id}',
        'io-impl' = 'org.apache.iceberg.rest.DlfFileIO'
    );

    Ganti {region-id} dengan ID wilayah katalog DLF Anda, misalnya cn-hangzhou atau ap-southeast-1. Untuk semua wilayah yang didukung beserta nilai titik akhirnya, lihat Wilayah dan titik akhir.

  5. Di pojok kanan bawah, klik Environment, pilih kluster sesi yang menjalankan VVR 11.2.0 atau versi lebih baru, lalu jalankan Pernyataan SQL.

Tabel berikut menjelaskan opsi konfigurasi.

Option Description Required Example
type Jenis katalog. Tetapkan ke iceberg. Yes iceberg
catalog-type Jenis katalog. Tetapkan ke rest. Yes rest
token.provider Penyedia token untuk otentikasi DLF. Tetapkan ke dlf. Yes dlf
uri Titik akhir Iceberg REST untuk katalog DLF Anda. Gunakan format http://{region-id}-vpc.dlf.aliyuncs.com/iceberg. Untuk nilai spesifik wilayah, lihat Wilayah dan titik akhir. Yes http://ap-southeast-1-vpc.dlf.aliyuncs.com/iceberg
warehouse Nama katalog DLF Anda. Yes iceberg_test
rest.signing-region ID wilayah katalog DLF Anda. Untuk daftar ID wilayah, lihat Wilayah dan titik akhir. Yes ap-southeast-1
io-impl Implementasi FileIO untuk DLF. Tetapkan ke org.apache.iceberg.rest.DlfFileIO. Yes org.apache.iceberg.rest.DlfFileIO

Konfigurasi Flink CDC untuk menggunakan katalog

Untuk membuat pekerjaan ingesti data, lihat Mengembangkan pekerjaan ingesti data Flink CDC.

Jika Anda telah membuat pemetaan Katalog Flink, lihat Gunakan kembali Katalog yang ada untuk mendapatkan informasi koneksi, lalu tambahkan konfigurasi Sink berikut ke pekerjaan ingesti data Anda:

sink:
  type: iceberg
  using.built-in-catalog: catalog_name

Contoh konfigurasi

Contoh berikut menunjukkan pola umum sinkronisasi data menggunakan pekerjaan YAML Flink CDC.

Sinkronkan seluruh database MySQL ke DLF

Pekerjaan ini menyinkronkan seluruh database MySQL ke DLF:

source:
  type: mysql
  name: MySQL Source
  hostname: ${mysql.hostname}
  port: ${mysql.port}
  username: ${mysql.username}
  password: ${mysql.password}
  tables: mysql_test.\.*
  server-id: 8601-8604
  # (Opsional) Sinkronkan tabel yang dibuat selama fase inkremental.
  scan.binlog.newly-added-table.enabled: true
  # (Opsional) Sinkronkan komentar tabel dan bidang.
  include-comments.enabled: true
  # (Opsional) Utamakan chunk tak terbatas untuk mencegah error OOM pada TaskManager.
  scan.incremental.snapshot.unbounded-chunk-first.enabled: true
  # (Opsional) Percepat pembacaan dengan hanya mengurai tabel yang cocok.
  scan.only.deserialize.captured.tables.changelog.enabled: true

sink:
  type: iceberg
  using.built-in-catalog: catalog_name
Catatan

Kami merekomendasikan parameter opsional berikut untuk sumber MySQL. Untuk detailnya, lihat MySQL.

  1. Parameter: scan.binlog.newly-added-table.enabled

    Fungsi: Menyinkronkan tabel yang dibuat selama fase inkremental.

  2. Parameter: include-comments.enabled

    Fungsi: Menyinkronkan komentar tabel dan bidang.

  3. Parameter: scan.incremental.snapshot.unbounded-chunk-first.enabled

    Fungsi: Mencegah error OOM pada TaskManager.

  4. Parameter: scan.only.deserialize.captured.tables.changelog.enabled: true

    Fungsi: Mempercepat pembacaan dengan hanya mengurai tabel yang cocok.

Tulis ke tabel DLF yang dipartisi

Tentukan kunci partisi menggunakan parameter partition-keys. Untuk detailnya, lihat Referensi pengembangan pekerjaan ingesti data Flink CDC. Contoh:

source:
  type: mysql
  name: MySQL Source
  hostname: ${mysql.hostname}
  port: ${mysql.port}
  username: ${mysql.username}
  password: ${mysql.password}
  tables: mysql_test.\.*
  server-id: 8601-8604
  # (Opsional) Sinkronkan tabel yang dibuat selama fase inkremental.
  scan.binlog.newly-added-table.enabled: true
  # (Opsional) Sinkronkan komentar tabel dan bidang.
  include-comments.enabled: true
  # (Opsional) Utamakan chunk tak terbatas untuk mencegah error OOM pada TaskManager.
  scan.incremental.snapshot.unbounded-chunk-first.enabled: true
  # (Opsional) Percepat pembacaan dengan hanya mengurai tabel yang cocok.
  scan.only.deserialize.captured.tables.changelog.enabled: true

sink:
  type: iceberg
  using.built-in-catalog: catalog_name

transform:
  - source-table: mysql_test.tbl1
    # (Opsional) Tetapkan kunci partisi.
    partition-keys: id,pt
  - source-table: mysql_test.tbl2
    partition-keys: id,pt

Tulis ke tabel DLF append-only

Pekerjaan ingesti data menangkap semua jenis event perubahan dari sumber. Untuk menerapkan soft delete dengan mengonversi operasi DELETE menjadi operasi INSERT di downstream, konfigurasikan pekerjaan seperti yang ditunjukkan di bawah. Untuk detailnya, lihat Referensi pengembangan pekerjaan ingesti data Flink CDC.

source:
  type: mysql
  name: MySQL Source
  hostname: ${mysql.hostname}
  port: ${mysql.port}
  username: ${mysql.username}
  password: ${mysql.password}
  tables: mysql_test.\.*
  server-id: 8601-8604
  # (Opsional) Sinkronkan tabel yang dibuat selama fase inkremental.
  scan.binlog.newly-added-table.enabled: true
  # (Opsional) Sinkronkan komentar tabel dan bidang.
  include-comments.enabled: true
  # (Opsional) Utamakan chunk tak terbatas untuk mencegah error OOM pada TaskManager.
  scan.incremental.snapshot.unbounded-chunk-first.enabled: true
  # (Opsional) Percepat pembacaan dengan hanya mengurai tabel yang cocok.
  scan.only.deserialize.captured.tables.changelog.enabled: true

sink:
  type: iceberg
  using.built-in-catalog: catalog_name

transform:
  - source-table: mysql_test.tbl1
    # (Opsional) Tetapkan kunci partisi.
    partition-keys: id,pt
    # (Opsional) Terapkan soft delete.
    projection: \*, __data_event_type__ AS op_type
    converter-after-transform: SOFT_DELETE
  - source-table: mysql_test.tbl2
    # (Opsional) Tetapkan kunci partisi.
    partition-keys: id,pt
    # (Opsional) Terapkan soft delete.
    projection: \*, __data_event_type__ AS op_type
    converter-after-transform: SOFT_DELETE
Catatan
  • Menambahkan __data_event_type__ ke proyeksi menulis jenis event perubahan sebagai bidang baru di tabel downstream. Menetapkan converter-after-transform ke SOFT_DELETE mengonversi operasi DELETE menjadi INSERT, sehingga tabel downstream mencatat semua event perubahan. Untuk detailnya, lihat Referensi pengembangan pekerjaan ingesti data Flink CDC.

Sinkronkan data CDC Kafka real-time ke DLF

Asumsikan topik inventory di Kafka menyimpan data perubahan untuk dua tabel—customers dan products—dalam format Debezium JSON. Pekerjaan ini menyinkronkan data dari tabel-tabel tersebut ke target DLF yang sesuai:

source:
  type: kafka
  name: Kafka Source
  properties.bootstrap.servers: ${kafka.bootstrap.servers}
  topic: inventory
  scan.startup.mode: earliest-offset
  value.format: debezium-json
  debezium-json.distributed-tables: true

sink:
  type: iceberg
  using.built-in-catalog: catalog_name

# Debezium JSON tidak memiliki informasi primary key. Tambahkan secara manual.
transform:
  - source-table: \.*.\.*
    projection: \*
    primary-keys: id
Catatan
  • Sumber Kafka mendukung format canal-json, debezium-json (default), dan json.

  • Saat menggunakan debezium-json, tambahkan primary key secara manual menggunakan aturan transformasi, karena pesan Debezium JSON tidak berisi informasi primary key:

    transform:
      - source-table: \.*.\.*
        projection: \*
        primary-keys: id
  • Jika data satu tabel tersebar di beberapa partisi, atau tabel dari berbagai partisi perlu digabungkan, tetapkan debezium-json.distributed-tables atau canal-json.distributed-tables ke true.

  • Sumber Kafka mendukung beberapa strategi inferensi skema. Tetapkan preferensi Anda menggunakan parameter schema.inference.strategy. Untuk detailnya, lihat Kafka.

Sinkronkan log Kafka real-time ke DLF

Jika kluster Kafka Anda menyimpan data dalam format JSON kustom, konfigurasikan pekerjaan YAML Flink CDC untuk menyinkronkannya ke DLF. Sistem menangani inferensi tipe data, inferensi skema, dan evolusi skema secara otomatis.

Asumsikan topik inventory menyimpan data untuk satu tabel log dalam format JSON. Pekerjaan ini menyinkronkan data ke tabel DLF yang sesuai:

source:
  type: kafka
  name: Kafka Source
  properties.bootstrap.servers: ${kafka.bootstrap.servers}
  topic: inventory
  scan.startup.mode: earliest-offset
  value.format: json
  # (Opsional) Ratakan secara rekursif kolom bersarang dalam data JSON.
  json.infer-schema.flatten-nested-columns.enable: true
  # (Opsional) Lewati 100 error parsing pertama. Pekerjaan gagal jika error melebihi 100.
  ingestion.ignore-errors: true
  ingestion.error-tolerance.max-count: 100

sink:
  type: iceberg
  using.built-in-catalog: catalog_name

# Tambahkan primary key ke tabel.
transform:
  - source-table: \.*.\.*
    projection: \*
    primary-keys: id

# Tulis semua data topik inventory ke test_database.inventory.
route:
  - source-table: inventory
    sink-table: test_database.inventory

pipeline:
  # (Opsional) Catat data kotor yang menyebabkan exception pemrosesan.
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger

Asumsikan topik inventory menyimpan data untuk beberapa tabel log dalam format JSON, dengan nama database dan tabel berada di bidang databaseName dan tableName. Pekerjaan ini menyinkronkan data dari tabel-tabel tersebut ke target DLF yang sesuai:

source:
  type: kafka
  name: Kafka Source
  properties.bootstrap.servers: ${kafka.bootstrap.servers}
  topic: inventory
  scan.startup.mode: earliest-offset
  value.format: json
  # (Opsional) Ratakan secara rekursif kolom bersarang dalam data JSON.
  json.infer-schema.flatten-nested-columns.enable: true
  # Gunakan nilai bidang databaseName dan tableName sebagai nama database dan tabel.
  json.decode.parser-table-id.fields: databaseName,tableName
  # (Opsional) Lewati 100 error parsing pertama. Pekerjaan gagal jika error melebihi 100.
  ingestion.ignore-errors: true
  ingestion.error-tolerance.max-count: 100

sink:
  type: iceberg
  using.built-in-catalog: catalog_name

# Tambahkan primary key ke tabel.
transform:
  - source-table: \.*.\.*
    projection: \*
    primary-keys: id

# Tulis ods.inventory, ods.customer, dan ods.user ke test_database.inventory, test_database.customer, dan test_database.user masing-masing.
route:
  - source-table: ods.inventory
    sink-table: test_database.inventory
  - source-table: ods.customer
    sink-table: test_database.customer
  - source-table: ods.user
    sink-table: test_database.user

pipeline:
  # (Opsional) Catat data kotor yang menyebabkan exception pemrosesan.
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger

Untuk detail tentang strategi inferensi dan evolusi skema untuk tabel sumber Kafka dalam format JSON, lihat Strategi parsing skema dan sinkronisasi perubahan.

Untuk detail konfigurasi pekerjaan yang komprehensif, lihat Referensi pengembangan pekerjaan ingesti data Flink CDC.