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:
-
Ruang kerja Flink yang sepenuhnya dikelola. Lihat Aktifkan Realtime Compute for Apache Flink.
-
Ruang kerja Flink dan DLF berada di wilayah yang sama.
-
Virtual Private Cloud (VPC) tempat ruang kerja Flink berada telah ditambahkan ke daftar putih DLF. Lihat Konfigurasi daftar putih VPC.
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.
-
Masuk ke Konsol Manajemen Realtime Compute for Apache Flink.
-
Pada kolom Actions ruang kerja Anda, klik Console.
-
Di panel navigasi kiri, klik Development > Scripts.
-
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, misalnyacn-hangzhouatauap-southeast-1. Untuk semua wilayah yang didukung beserta nilai titik akhirnya, lihat Wilayah dan titik akhir. -
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
Kami merekomendasikan parameter opsional berikut untuk sumber MySQL. Untuk detailnya, lihat MySQL.
-
Parameter: scan.binlog.newly-added-table.enabled
Fungsi: Menyinkronkan tabel yang dibuat selama fase inkremental.
-
Parameter: include-comments.enabled
Fungsi: Menyinkronkan komentar tabel dan bidang.
-
Parameter: scan.incremental.snapshot.unbounded-chunk-first.enabled
Fungsi: Mencegah error OOM pada TaskManager.
-
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
-
Menambahkan
__data_event_type__ke proyeksi menulis jenis event perubahan sebagai bidang baru di tabel downstream. Menetapkanconverter-after-transformkeSOFT_DELETEmengonversi 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
-
Sumber Kafka mendukung format
canal-json,debezium-json(default), danjson. -
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-tablesataucanal-json.distributed-tablesketrue. -
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.