Konektor Kafka dapat digunakan sebagai sumber atau sink dalam pekerjaan ingesti data Flink CDC. Topik ini menjelaskan sintaksis, parameter, dan contoh penggunaan konektor Kafka YAML.
Prasyarat
Hubungkan ke kluster Anda dengan salah satu metode berikut:
-
Hubungkan ke kluster ApsaraMQ for Kafka
-
Kluster Kafka menjalankan versi 0.11 atau yang lebih baru.
-
Kluster ApsaraMQ for Kafka telah dibuat. Untuk informasi selengkapnya, lihat Buat sumber daya.
-
Ruang kerja Flink dan kluster Kafka berada dalam VPC yang sama, dan ApsaraMQ for Kafka telah menambahkan Flink ke daftar putihnya. Untuk informasi selengkapnya, lihat Konfigurasi daftar putih.
PentingBatasan penulisan ke ApsaraMQ for Kafka:
-
ApsaraMQ for Kafka tidak mendukung penulisan data dalam format kompresi zstd.
-
ApsaraMQ for Kafka tidak mendukung penulisan idempoten atau penulisan transaksional, sehingga semantik tepat-sekali yang disediakan oleh tabel sink Kafka tidak tersedia. Mulai dari VVR 8.0.0, Klien Kafka open source yang digunakan oleh Konektor Kafka ditingkatkan ke versi 3.x, di mana properti
properties.enable.idempotencesecara default bernilai true, yang secara eksplisit mengaktifkan penulisan idempoten. Oleh karena itu, saat Anda menulis ke ApsaraMQ for Kafka pada VVR 8.0.0 atau yang lebih baru, Anda harus secara eksplisit menambahkanproperties.enable.idempotence=falseke tabel sink untuk menonaktifkan penulisan idempoten dan menghindari kegagalan penulisan. Untuk perbandingan mesin penyimpanan dan batasan fitur ApsaraMQ for Kafka, lihat Perbandingan mesin penyimpanan.
-
-
Hubungkan ke kluster Apache Kafka yang dikelola sendiri
-
Kluster Apache Kafka yang dikelola sendiri menjalankan versi 0.11 atau yang lebih baru.
-
Jaringan antara Flink dan kluster Apache Kafka yang dikelola sendiri terhubung. Untuk menghubungkan ke kluster yang dikelola sendiri melalui Internet, lihat Koneksi jaringan.
-
Hanya item konfigurasi klien Apache Kafka 2.8 yang didukung. Untuk informasi selengkapnya, lihat dokumentasi konfigurasi konsumen dan produsen Apache Kafka.
-
Batasan
-
Kami merekomendasikan agar Anda menggunakan Kafka sebagai sumber data untuk ingesti data Flink CDC pada VVR 11.1 atau yang lebih baru.
-
Hanya format JSON, Debezium JSON, dan Canal JSON yang didukung. Format data lain tidak didukung.
-
Untuk sumber, data dari tabel yang sama yang didistribusikan di beberapa partisi hanya didukung pada VVR 8.0.11 dan yang lebih baru.
Catatan penggunaan
Saat ini, penulisan transaksional tidak direkomendasikan karena keterbatasan desain di komunitas Flink dan Kafka. Saat Anda mengatur sink.delivery-guarantee = exactly-once, Konektor Kafka mengaktifkan penulisan transaksional. Tiga masalah berikut diketahui:
-
Setiap Checkpoint menghasilkan ID transaksi. Jika interval Checkpoint terlalu pendek, terlalu banyak ID transaksi yang dihasilkan. Koordinator kluster Kafka mungkin kehabisan memori, yang membahayakan stabilitas kluster Kafka.
-
Setiap transaksi membuat instans Produsen. Jika terlalu banyak transaksi dikomit secara bersamaan, Pengelola Tugas (TaskManager) mungkin kehabisan memori, yang membahayakan stabilitas pekerjaan Flink.
-
Jika beberapa pekerjaan Flink menggunakan
sink.transactional-id-prefixyang sama, ID transaksi yang dihasilkan mungkin bentrok. Saat satu pekerjaan gagal menulis data, LSO (Log Start Offset) partisi Kafka diblokir sehingga tidak bisa maju, yang memengaruhi semua konsumen yang membaca data dari partisi tersebut.
Jika Anda memerlukan semantik tepat-sekali, gunakan konektor Upsert Kafka untuk menulis data ke tabel kunci primer dan andalkan kunci primer untuk memastikan idempotensi. Jika Anda harus menggunakan penulisan transaksional, lihat Semantik tepat-sekali.
Sintaksis
source:
type: kafka
name: Kafka source
properties.bootstrap.servers: localhost:9092
topic: ${kafka.topic}
sink:
type: kafka
name: Kafka Sink
properties.bootstrap.servers: localhost:9092
Item konfigurasi
-
Umum
Parameter
Deskripsi
Wajib
Tipe data
Nilai default
Catatan
type
Jenis sumber atau sink.
Ya
String
None
Nilainya harus kafka.
name
Nama sumber atau sink.
Tidak
String
None
None
properties.bootstrap.servers
Alamat broker Kafka.
Ya
String
None
Formatnya adalah
host:port,host:port,host:port. Gunakan koma (,) untuk memisahkan beberapa alamat.properties.*
Konfigurasi yang diteruskan langsung ke klien Kafka.
Tidak
String
None
Akhiran harus berupa item konfigurasi produsen atau konsumen yang ditentukan dalam dokumentasi resmi Apache Kafka.
Flink menghapus awalan properties. dan meneruskan konfigurasi yang tersisa ke klien Kafka. Misalnya, Anda dapat menggunakan
'properties.allow.auto.create.topics' = 'false'untuk menonaktifkan pembuatan topik otomatis.key.format
Format yang digunakan untuk membaca atau menulis bagian kunci pesan Kafka.
Tidak
String
None
-
Untuk sumber, hanya json yang didukung.
-
Untuk sink, nilai yang valid:
-
csv
-
json
-
CatatanParameter ini hanya didukung pada VVR 11.0.0 dan yang lebih baru.
value.format
Format yang digunakan untuk membaca atau menulis bagian nilai pesan Kafka.
Tidak
String
debezium-json
-
Untuk sumber, nilai yang valid:
-
debezium-json
-
canal-json
-
json
-
-
Untuk sink, nilai yang valid:
-
debezium-json
-
canal-json
-
canal-protobuf
-
Catatan-
Format debezium-json dan canal-json hanya didukung pada VVR 8.0.10 dan yang lebih baru.
-
Format json hanya didukung pada VVR 11.0.0 dan yang lebih baru.
-
-
Sumber
Parameter
Deskripsi
Wajib
Tipe data
Nilai default
Catatan
topic
Nama topik yang akan dibaca.
Tidak
String
None
Pisahkan beberapa nama topik dengan titik koma (;), misalnya, topic-1 dan topic-2.
CatatanAnda hanya dapat menentukan salah satu dari topic dan topic-pattern.
topic-pattern
Ekspresi reguler yang digunakan untuk mencocokkan nama topik yang akan dibaca. Semua topik yang cocok dengan ekspresi reguler akan dibaca saat pekerjaan dijalankan.
Tidak
String
None
Contoh:
-
user_event_.*: mencocokkan semua topik yang namanya dimulai denganuser_event_ -
prod\.logs\..*: mencocokkan topik yang namanya memiliki awalanprod.logs.(.harus di-escape)
CatatanAnda hanya dapat menentukan salah satu dari topic dan topic-pattern.
properties.group.id
ID kelompok konsumen.
Tidak
String
None
Jika ID grup yang ditentukan digunakan untuk pertama kalinya, Anda harus mengatur properties.auto.offset.reset ke earliest atau latest untuk menentukan offset startup awal.
scan.startup.mode
Offset awal tempat Kafka membaca data.
Tidak
String
group-offsets
Nilai yang valid:
-
earliest-offset: membaca data dari offset paling awal partisi Kafka.
-
latest-offset: membaca data dari offset paling akhir partisi Kafka.
-
group-offsets (default): membaca data dari offset yang dikomit oleh kelompok konsumen yang ditentukan oleh properties.group.id.
-
timestamp: membaca data dari timestamp yang ditentukan oleh scan.startup.timestamp-millis.
-
specific-offsets: membaca data dari offset yang ditentukan oleh scan.startup.specific-offsets.
CatatanParameter ini hanya berlaku saat pekerjaan dimulai tanpa state. Saat pekerjaan dimulai ulang dari Checkpoint atau dipulihkan dari state, pembacaan dilanjutkan dari kemajuan yang disimpan dalam state.
scan.startup.specific-offsets
Offset startup setiap partisi saat mode startup adalah specific-offsets.
Tidak
String
None
Contoh:
partition:0,offset:42;partition:1,offset:300scan.startup.timestamp-millis
Timestamp startup saat mode startup adalah timestamp.
Tidak
Long
None
Unit: milidetik.
scan.topic-partition-discovery.interval
Interval deteksi dinamis topik dan partisi Kafka.
Tidak
Duration
5 menit
Interval penemuan partisi default adalah 5 menit. Untuk menonaktifkan fitur ini, atur secara eksplisit interval penemuan partisi ke nilai non-positif. Saat penemuan partisi dinamis diaktifkan, sumber Kafka secara otomatis menemukan partisi baru dan membaca data darinya. Dalam mode topic-pattern, sumber membaca data dari partisi baru topik yang ada dan dari semua partisi topik baru yang cocok dengan ekspresi reguler.
scan.check.duplicated.group.id
Menentukan apakah akan memeriksa apakah kelompok konsumen yang ditentukan oleh
properties.group.iddiduplikasi.Tidak
Boolean
false
Nilai yang valid:
-
true: memeriksa konflik kelompok konsumen sebelum pekerjaan dimulai. Jika terjadi konflik, pekerjaan melaporkan error. Ini menghindari konflik dengan kelompok konsumen yang sudah ada.
-
false: memulai pekerjaan langsung tanpa memeriksa konflik kelompok konsumen.
schema.inference.strategy
Kebijakan penguraian Schema.
Tidak
String
continuous
Nilai yang valid:
-
continuous: mengurai Schema setiap record. Saat dua Schema berturut-turut tidak kompatibel, Schema yang lebih luas diurai dan event perubahan Schema dihasilkan.
-
static: mengurai Schema hanya sekali saat pekerjaan dimulai. Data kemudian diurai berdasarkan Schema awal, dan tidak ada event perubahan Schema yang dihasilkan.
Catatan-
Untuk informasi selengkapnya tentang penguraian Schema, lihat Kebijakan penguraian dan sinkronisasi perubahan Schema.
-
Parameter ini hanya didukung pada VVR 8.0.11 dan yang lebih baru.
scan.max.pre.fetch.records
Jumlah maksimum pesan yang dikonsumsi dan diurai dari setiap partisi selama penguraian Schema awal.
Tidak
Int
50
Sebelum pekerjaan membaca dan memproses data, jumlah pesan terbaru yang ditentukan dikonsumsi terlebih dahulu dari setiap partisi untuk menginisialisasi informasi Schema.
key.fields-prefix
Awalan kustom yang ditambahkan ke nama field yang diurai dari kunci pesan. Ini menghindari konflik penamaan setelah kunci pesan Kafka diurai.
Tidak
String
None
Misalnya, jika parameter ini diatur ke key_ dan kunci berisi field bernama a, field tersebut akan diberi nama key_a setelah kunci diurai.
CatatanNilai key.fields-prefix tidak boleh menjadi awalan dari value.fields-prefix.
value.fields-prefix
Awalan kustom yang ditambahkan ke nama field yang diurai dari nilai pesan. Ini menghindari konflik penamaan setelah nilai pesan Kafka diurai.
Tidak
String
None
Misalnya, jika parameter ini diatur ke value_ dan nilai berisi field bernama b, field tersebut akan diberi nama value_b setelah nilai diurai.
CatatanNilai value.fields-prefix tidak boleh menjadi awalan dari key.fields-prefix.
metadata.list
Kolom metadata yang diteruskan ke downstream.
Tidak
String
None
Kolom metadata yang tersedia adalah
topic,partition,offset,timestamp,timestamp-type,headers,leader-epoch,__raw_key__, dan__raw_value__. Pisahkan beberapa kolom dengan koma (,).CatatanKolom metadata
__raw_key__dan__raw_value__hanya tersedia pada VVR 11.6 dan yang lebih baru.scan.value.initial-schemas.ddls
Schema awal tabel tertentu, ditentukan dengan menggunakan pernyataan DDL.
Tidak
String
None
Pisahkan beberapa pernyataan DDL dengan
;. Misalnya, gunakanCREATE TABLE db1.t1 (id BIGINT, name VARCHAR(10)); CREATE TABLE db1.t2 (id BIGINT);untuk menentukan Schema awal tabel db1.t1 dan db1.t2.Schema tabel dalam pernyataan DDL harus konsisten dengan tabel tujuan dan mematuhi aturan sintaksis Flink SQL.
CatatanParameter ini hanya didukung pada VVR 11.5 dan yang lebih baru.
ingestion.ignore-errors
Menentukan apakah akan mengabaikan error selama penguraian data.
Tidak
Boolean
false
CatatanParameter ini hanya didukung pada VVR 11.5 dan yang lebih baru.
ingestion.error-tolerance.max-count
Jumlah akumulasi error penguraian setelah pekerjaan gagal saat error penguraian data diabaikan.
Tidak
Integer
-1
Parameter ini hanya berlaku saat ingestion.ignore-errors diaktifkan. Nilai default -1 menunjukkan bahwa exception penguraian tidak menyebabkan pekerjaan gagal.
CatatanParameter ini hanya didukung pada VVR 11.5 dan yang lebih baru.
scan.duplicate-field.strategy
Menentukan cara menangani nama field duplikat yang diurai dari bagian kunci dan nilai.
Tidak
String
EXCEPTION
Nilai yang valid:
-
EXCEPTION: melempar exception saat field duplikat ada di kunci dan nilai. Ini adalah perilaku default pada VVR 11.6 dan yang lebih awal.
-
PREFER_KEY: lebih memilih nilai field kunci saat field diduplikasi.
-
PREFER_VALUE: lebih memilih nilai field nilai saat field diduplikasi.
CatatanParameter ini hanya didukung pada VVR 11.7 dan yang lebih baru.
-
Tabel sumber dalam format Debezium JSON
Parameter
Wajib
Tipe data
Nilai default
Deskripsi
debezium-json.distributed-tables
Tidak
Boolean
false
Aktifkan opsi ini jika data satu tabel dalam Debezium JSON muncul di beberapa partisi.
CatatanParameter ini hanya didukung pada VVR 8.0.11 dan yang lebih baru.
PentingSetelah Anda mengubah parameter ini, Anda harus memulai pekerjaan tanpa state.
debezium-json.schema-include
Tidak
Boolean
false
Saat Anda mengonfigurasi Debezium Kafka Connect, Anda dapat mengaktifkan konfigurasi Kafka value.converter.schemas.enable untuk menyertakan schema dalam pesan. Opsi ini menentukan apakah pesan Debezium JSON berisi schema.
Nilai yang valid:
-
true: pesan Debezium JSON berisi schema.
-
false: pesan Debezium JSON tidak berisi schema.
debezium-json.ignore-parse-errors
Tidak
Boolean
false
Nilai yang valid:
-
true: melewati baris saat ini saat terjadi exception penguraian.
-
false (default): melaporkan error, dan pekerjaan gagal dimulai.
debezium-json.infer-schema.primitive-as-string
Tidak
Boolean
false
Menentukan apakah akan mengurai semua tipe sebagai String saat Schema tabel diurai.
Nilai yang valid:
-
true: mengurai semua tipe primitif sebagai String.
-
false (default): mengurai data berdasarkan aturan dasar.
debezium-json.infer-schema.string-type-inference
Tidak
Boolean
true
Menentukan apakah akan mencoba menginfer field string sebagai tipe TIME, DATE, atau TIMESTAMP. Jika parameter ini diatur ke false, inferensi dilewati dan field tetap STRING.
CatatanParameter ini hanya didukung pada VVR 11.8 dan yang lebih baru.
-
-
Tabel sumber dalam format Canal JSON
Parameter
Wajib
Tipe data
Nilai default
Deskripsi
canal-json.distributed-tables
Tidak
Boolean
false
Aktifkan opsi ini jika data satu tabel dalam Canal JSON muncul di beberapa partisi.
CatatanParameter ini hanya didukung pada VVR 8.0.11 dan yang lebih baru.
PentingSetelah Anda mengubah parameter ini, Anda harus memulai pekerjaan tanpa state.
canal-json.database.include
Tidak
String
None
Ekspresi reguler opsional yang mencocokkan field metadata database dalam record Canal sehingga hanya record changelog database yang ditentukan yang dibaca. String ekspresi reguler kompatibel dengan Java Pattern.
canal-json.table.include
Tidak
String
None
Ekspresi reguler opsional yang mencocokkan field metadata tabel dalam record Canal sehingga hanya record changelog tabel yang ditentukan yang dibaca. String ekspresi reguler kompatibel dengan Java Pattern.
canal-json.ignore-parse-errors
Tidak
Boolean
false
Nilai yang valid:
-
true: melewati baris saat ini saat terjadi exception penguraian.
-
false (default): melaporkan error, dan pekerjaan gagal dimulai.
canal-json.infer-schema.primitive-as-string
Tidak
Boolean
false
Menentukan apakah akan mengurai semua tipe sebagai String saat Schema tabel diurai.
Nilai yang valid:
-
true: mengurai semua tipe primitif sebagai String.
-
false (default): mengurai data berdasarkan aturan dasar.
canal-json.infer-schema.strategy
Tidak
String
AUTO
Kebijakan penguraian yang digunakan saat Schema tabel diurai.
Nilai yang valid:
-
AUTO (default): mengurai tipe secara otomatis dengan mengurai data JSON. Jika data tidak berisi field sqlType, kami merekomendasikan Anda menggunakan AUTO untuk menghindari kegagalan penguraian.
-
SQL_TYPE: mengurai tipe berdasarkan array sqlType dalam data Canal JSON. Jika data berisi field sqlType, kami merekomendasikan Anda mengatur canal-json.infer-schema.strategy ke SQL_TYPE untuk tipe yang lebih tepat.
-
MYSQL_TYPE: mengurai tipe berdasarkan array mysqlType dalam data Canal JSON.
Saat data Canal JSON di Kafka berisi field sqlType dan Anda memerlukan pemetaan tipe yang lebih tepat, kami merekomendasikan Anda mengatur canal-json.infer-schema.strategy ke SQL_TYPE.
Untuk aturan pemetaan tipe sqlType, lihat Penguraian schema Canal JSON.
Catatan-
Parameter ini hanya didukung pada VVR 11.1 dan yang lebih baru.
-
MYSQL_TYPE hanya didukung pada VVR 11.3 dan yang lebih baru.
canal-json.mysql.treat-mysql-timestamp-as-datetime-enabled
Tidak
Boolean
true
Menentukan apakah akan memetakan tipe timestamp MySQL ke tipe timestamp CDC:
-
true (default): memetakan tipe timestamp MySQL ke tipe timestamp CDC.
-
false: memetakan tipe timestamp MySQL ke tipe timestamp_ltz CDC.
canal-json.mysql.treat-tinyint1-as-boolean.enabled
Tidak
Boolean
true
Menentukan apakah akan memetakan tipe tinyint(1) MySQL ke tipe boolean CDC saat penguraian MYSQL_TYPE digunakan:
-
true (default): memetakan tipe tinyint(1) MySQL ke tipe boolean CDC.
-
false: memetakan tipe tinyint(1) MySQL ke tipe tinyint(1) CDC.
Parameter ini hanya berlaku saat canal-json.infer-schema.strategy diatur ke MYSQL_TYPE.
canal-json.infer-schema.string-type-inference
Tidak
Boolean
true
Menentukan apakah akan mencoba menginfer field string sebagai tipe TIME, DATE, atau TIMESTAMP. Jika parameter ini diatur ke false, inferensi dilewati dan field tetap STRING.
CatatanParameter ini hanya didukung pada VVR 11.8 dan yang lebih baru.
-
-
Tabel sumber dalam format JSON
Parameter
Wajib
Tipe data
Nilai default
Deskripsi
json.timestamp-format.standard
Tidak
String
SQL
Format timestamp input dan output. Nilai yang valid:
-
SQL: mengurai timestamp input dalam format yyyy-MM-dd HH:mm:ss.s{precision}, misalnya, 2020-12-30 12:13:14.123.
-
ISO-8601: mengurai timestamp input dalam format yyyy-MM-ddTHH:mm:ss.s{precision}, misalnya, 2020-12-30T12:13:14.123.
json.ignore-parse-errors
Tidak
Boolean
false
Nilai yang valid:
-
true: melewati baris saat ini saat terjadi exception penguraian.
-
false (default): melaporkan error, dan pekerjaan gagal dimulai.
json.infer-schema.primitive-as-string
Tidak
Boolean
false
Menentukan apakah akan mengurai semua tipe sebagai String saat Schema tabel diurai.
Nilai yang valid:
-
true: mengurai semua tipe primitif sebagai String.
-
false (default): mengurai data berdasarkan aturan dasar.
json.infer-schema.flatten-nested-columns.enable
Tidak
Boolean
false
Menentukan apakah akan memperluas kolom bersarang secara rekursif dalam data JSON selama penguraian. Nilai yang valid:
-
true: memperluas kolom bersarang secara rekursif.
-
false (default): memperlakukan kolom bersarang sebagai String.
json.decode.parser-table-id.fields
Tidak
String
None
Menentukan apakah akan menggunakan nilai field JSON tertentu untuk menghasilkan tableId saat data JSON diurai. Pisahkan beberapa field dengan
,. Misalnya, jika data JSON adalah{"col0":"a", "col1","b", "col2","c"}, hasil yang dihasilkan adalah sebagai berikut:Konfigurasi
tableId
col0
a
col0,col1
a.b
col0,col1,col2
a.b.c
json.infer-schema.fixed-types
Tidak
String
None
Tipe tertentu dari field tertentu saat data JSON diurai. Pisahkan beberapa field dengan
,. Misalnya,id BIGINT, name VARCHAR(10)menentukan tipe field id dalam data JSON sebagai BIGINT dan tipe field name sebagai VARCHAR(10).Catatan-
Parameter ini hanya didukung pada VVR 11.5 dan yang lebih baru.
-
Saat Anda menggunakan parameter ini pada VVR 11.5, Anda juga harus menambahkan
scan.max.pre.fetch.records: 0.
json.decode.converter-class
Tidak
String
None
Nama lengkap kelas implementasi. Konverter dapat memodifikasi byte[] JSON sebelum data JSON diurai.
CatatanParameter ini hanya didukung pada VVR 11.6 dan yang lebih baru.
json.decode.empty-value-as-delete.enabled
Tidak
Boolean
false
Menentukan apakah akan mengurai pesan tombstone (dengan nilai kosong) dalam topik Kafka yang dikompaksi sebagai event DELETE. Ini berlaku untuk skenario di mana nilai kosong menunjukkan penghapusan, seperti pencerminkan topik yang dikompaksi dan sinyal penghapusan CDC.
CatatanParameter ini hanya didukung pada VVR 11.7 dan yang lebih baru.
json.infer-schema.string-type-inference
Tidak
Boolean
true
Menentukan apakah akan mencoba menginfer field string sebagai tipe TIME, DATE, atau TIMESTAMP. Jika parameter ini diatur ke false, inferensi dilewati dan field tetap STRING.
CatatanParameter ini hanya didukung pada VVR 11.8 dan yang lebih baru.
-
-
-
Sink
Parameter
Deskripsi
Wajib
Tipe data
Nilai default
Catatan
type
Jenis sink.
Ya
String
None
Nilainya harus kafka.
name
Nama sink.
Tidak
String
None
None
topic
Nama topik Kafka.
Tidak
String
None
Jika opsi ini diaktifkan, semua data ditulis ke topik ini.
CatatanJika opsi ini dinonaktifkan, setiap record ditulis ke topik yang namanya sesuai dengan string TableID-nya (dihasilkan dengan menggabungkan menggunakan
.), misalnya,databaseName.tableName.partition.strategy
Kebijakan yang digunakan untuk menulis data ke partisi Kafka.
Tidak
String
all-to-zero
Nilai yang valid:
-
all-to-zero (default): menulis semua data ke partisi 0.
-
hash-by-key: menulis data ke beberapa partisi berdasarkan nilai hash kunci primer. Ini memastikan bahwa data dengan kunci primer yang sama ditulis ke partisi yang sama secara berurutan.
sink.tableId-to-topic.mapping
Pemetaan antara nama tabel sumber dan nama topik Kafka tujuan.
Tidak
String
None
Pisahkan setiap pemetaan dengan
;. Pisahkan nama tabel sumber dan nama topik Kafka tujuan dengan:. Nama tabel dapat berupa ekspresi reguler. Beberapa tabel yang dipetakan ke topik yang sama dapat digabungkan dengan,. Misalnya:mydb.mytable1:topic1;mydb.mytable2:topic2.CatatanParameter ini memungkinkan Anda mengubah topik yang dipetakan sambil mempertahankan informasi nama tabel asli.
-
Tabel sink dalam format Canal JSON
Parameter
Wajib
Tipe data
Nilai default
Deskripsi
canal-json.serialize.update.keep-changed-fields-only
Tidak
Boolean
false
Menentukan apakah bagian lama pesan UPDATE dalam format canal-json hanya mencakup nilai pra-perubahan dari field yang berubah.
CatatanParameter ini hanya didukung pada VVR 11.8 dan versi yang lebih baru.
-
Tabel sink dalam format Debezium JSON
Parameter
Wajib
Tipe data
Nilai default
Deskripsi
debezium-json.include-schema.enabled
Tidak
Boolean
false
Menentukan apakah data Debezium JSON mencakup informasi schema.
debezium-json.emit.full-table-id.enabled
Tidak
Boolean
false
Menentukan apakah ID tabel tiga bagian lengkap dituliskan ke field metadata Debezium JSON.
Pemetaan saat parameter ini diaktifkan:
Bagian ID Tabel CDC
Debezium JSON Key
Namespace
dbSchema
schemaTabel
tablePemetaan saat parameter ini dinonaktifkan:
Bagian ID Tabel CDC
Kunci Debezium JSON
Namespace
None
Schema
dbTabel
tableCatatanParameter ini hanya didukung pada VVR 11.6 dan versi yang lebih baru.
-
Gunakan katalog yang sudah ada
Mulai dari VVR 11.5, Anda dapat langsung mereferensikan Katalog Kafka bawaan yang dibuat di halaman Katalog dalam pekerjaan ingesti data Flink CDC. Ini mengurangi kebutuhan untuk menulis properti koneksi secara manual.
source:
type: kafka
using.built-in-catalog: kafka_catalog
Saat ini, pekerjaan ingesti data mendukung penggunaan ulang otomatis parameter Katalog Kafka berikut:
-
properties.bootstrap.servers
-
format
-
key.fields-prefix
-
value.fields-prefix
-
timestamp-format.standard
-
infer-schema.flatten-nested-columns.enable
-
infer-schema.primitive-as-string
-
max.fetch.records
Untuk mengganti parameter yang digunakan ulang secara otomatis, tentukan secara eksplisit parameter YAML yang sesuai, yang memiliki prioritas lebih tinggi.
Contoh konfigurasi
-
Gunakan Kafka sebagai sumber pekerjaan ingesti data:
source: type: kafka name: Kafka source properties.bootstrap.servers: ${kafka.bootstraps.server} topic: ${kafka.topic} value.format: ${value.format} scan.startup.mode: ${scan.startup.mode} sink: type: hologres name: Hologres sink endpoint: <yourEndpoint> dbname: <yourDbname> username: ${secret_values.ak_id} password: ${secret_values.ak_secret} sink.type-normalize-strategy: BROADEN -
Gunakan Kafka sebagai sink pekerjaan ingesti data:
source: type: mysql name: MySQL Source hostname: ${secret_values.mysql.hostname} port: ${mysql.port} username: ${secret_values.mysql.username} password: ${secret_values.mysql.password} tables: ${mysql.source.table} server-id: 8601-8604 sink: type: kafka name: Kafka Sink properties.bootstrap.servers: ${kafka.bootstraps.server} route: - source-table: ${mysql.source.table} sink-table: ${kafka.topic}Dalam contoh ini, modul route digunakan untuk mengatur nama topik Kafka tempat tabel sumber ditulis.
ApsaraMQ for Kafka tidak mengaktifkan pembuatan topik otomatis secara default. Untuk informasi selengkapnya, lihat FAQ tentang pembuatan topik otomatis. Sebelum Anda menulis data ke ApsaraMQ for Kafka, Anda harus membuat topik yang sesuai. Untuk informasi selengkapnya, lihat Langkah 3: Buat sumber daya.
Contoh
Contoh berikut menunjukkan konfigurasi untuk skenario umum.
Baca satu topik
Contoh berikut membaca topik customers dan menulis data ke Data Lake Formation:
source:
type: kafka
topic: customers
properties.bootstrap.servers: localhost:9092
properties.group.id: ${kafka.group.id}
value.format: json
# (可选)动态识别每条数据Schema,并比对生成Schema变更
schema.inference.strategy: continuous
sink:
type: paimon
name: Paimon Sink
catalog.properties.metastore: rest
catalog.properties.uri: dlf_uri
catalog.properties.warehouse: your_warehouse
catalog.properties.token.provider: dlf
#(可选)提交用户名,建议为不同作业设置不同的提交用户以避免冲突
commit.user: your_job_name
#(可选)开启删除向量,提升读取性能
table.properties.deletion-vectors.enabled: true
Dalam contoh ini, nama tabel yang dihasilkan untuk format JSON sama dengan nama topik secara default.
Baca beberapa topik
Contoh berikut membaca beberapa topik yang cocok dengan ekspresi reguler dan menulis data ke konektor StarRocks:
source:
type: kafka
topic-pattern: user_event_.*
properties.bootstrap.servers: localhost:9092
properties.group.id: ${kafka.group.id}
value.format: json
# (可选)动态识别每条数据Schema,并比对生成Schema变更
schema.inference.strategy: continuous
sink:
type: starrocks
jdbc-url: jdbc:mysql://<yourFeHostname>:9030
load-url: <yourFeHostname>:8030
username: <yourUsername>
password: ${secret_values.starrocks_password}
# (可选)数据量不大的作业建议调低 flush 间隔,避免数据长时间不落盘(默认 300000,即 5 分钟)
sink.buffer-flush.interval-ms: 5000
# (可选)上游为 utf8mb4 字符集时建议设为 4,避免文本截断(默认 3)
unicode-char.max-bytes: 4
# (可选)自动建表的分桶数;StarRocks 2.5.7 以下必须显式配置,更高版本可自动推断
table.create.num-buckets: 8
# (可选)自动建表的副本数,按集群情况配置
table.create.properties.replication_num: 3
# (可选)StarRocks 3.2 及以上建议开启,加速表结构变更
table.create.properties.fast_schema_evolution: true
# 注意:通过 transform 变更主键时,必须同时设置 sink.ignore.update-before: false,
# 否则旧主键对应的行会残留在下游
Dalam contoh ini, nama tabel yang dihasilkan untuk format JSON sama dengan nama topik secara default.
Baca kunci dan cegah konflik field
Untuk mencegah error yang disebabkan oleh konflik field antara kunci dan nilai, gunakan salah satu solusi berikut:
-
Tambahkan awalan ke field untuk menghindari konflik field:
source:
type: kafka
topic: ${kafka.topic}
properties.bootstrap.servers: localhost:9092
properties.group.id: ${kafka.group.id}
key.format: json
value.format: json
# key部分的字段名添加key_前缀
key.fields-prefix: key_
# value部分的字段名添加value_前缀
value.fields-prefix: value_
# (可选)动态识别每条数据Schema,并比对生成Schema变更
schema.inference.strategy: continuous
sink:
type: paimon
name: Paimon Sink
catalog.properties.metastore: rest
catalog.properties.uri: dlf_uri
catalog.properties.warehouse: your_warehouse
catalog.properties.token.provider: dlf
#(可选)提交用户名,建议为不同作业设置不同的提交用户以避免冲突
commit.user: your_job_name
#(可选)开启删除向量,提升读取性能
table.properties.deletion-vectors.enabled: true
-
Konfigurasikan kebijakan resolusi konflik. Untuk informasi selengkapnya, lihat scan.duplicate-field.strategy. Dalam konfigurasi berikut, field di kunci didahulukan, dan field dengan nama yang sama di nilai diabaikan saat terjadi duplikasi.
source:
type: kafka
topic: ${kafka.topic}
properties.bootstrap.servers: localhost:9092
properties.group.id: ${kafka.group.id}
key.format: json
value.format: json
# 优先使用key的字段,忽略同名的value中字段
scan.duplicate-field.strategy: PREFER_KEY
# (可选)动态识别每条数据Schema,并比对生成Schema变更
schema.inference.strategy: continuous
sink:
type: paimon
name: Paimon Sink
catalog.properties.metastore: rest
catalog.properties.uri: dlf_uri
catalog.properties.warehouse: your_warehouse
catalog.properties.token.provider: dlf
#(可选)提交用户名,建议为不同作业设置不同的提交用户以避免冲突
commit.user: your_job_name
#(可选)开启删除向量,提升读取性能
table.properties.deletion-vectors.enabled: true
Tambahkan kolom metadata
Contoh berikut membaca topik customers dan menambahkan kolom metadata topic dan partition ke field:
source:
type: kafka
topic: customers
properties.bootstrap.servers: localhost:9092
properties.group.id: ${kafka.group.id}
value.format: json
# (可选)动态识别每条数据Schema,并比对生成Schema变更
schema.inference.strategy: continuous
metadata.list: topic,partition
sink:
type: paimon
name: Paimon Sink
catalog.properties.metastore: rest
catalog.properties.uri: dlf_uri
catalog.properties.warehouse: your_warehouse
catalog.properties.token.provider: dlf
#(可选)提交用户名,建议为不同作业设置不同的提交用户以避免冲突
commit.user: your_job_name
#(可选)开启删除向量,提升读取性能
table.properties.deletion-vectors.enabled: true
Tangani error penguraian
Error penguraian data menyebabkan kegagalan pekerjaan. Anda dapat mengonfigurasi pekerjaan untuk mentoleransi error penguraian, yang biasanya digunakan bersama dengan Pengumpulan data kotor.
Contoh berikut sepenuhnya mengabaikan error penguraian:
source:
type: kafka
topic: customers
properties.bootstrap.servers: localhost:9092
properties.group.id: ${kafka.group.id}
value.format: json
# (可选)动态识别每条数据Schema,并比对生成Schema变更
schema.inference.strategy: continuous
# 开启忽略解析报错,默认忽略全部解析报错
ingestion.ignore-errors: true
sink:
type: paimon
name: Paimon Sink
catalog.properties.metastore: rest
catalog.properties.uri: dlf_uri
catalog.properties.warehouse: your_warehouse
catalog.properties.token.provider: dlf
#(可选)提交用户名,建议为不同作业设置不同的提交用户以避免冲突
commit.user: your_job_name
#(可选)开启删除向量,提升读取性能
table.properties.deletion-vectors.enabled: true
# 开启脏数据收集器,打印解析失败数据
pipeline:
dirty-data.collector:
name: Logger Dirty Data Collector
type: logger
Contoh berikut membuat pekerjaan gagal setelah 30 kegagalan penguraian terakumulasi:
source:
type: kafka
topic: customers
properties.bootstrap.servers: localhost:9092
properties.group.id: ${kafka.group.id}
value.format: json
# (可选)动态识别每条数据Schema,并比对生成Schema变更
schema.inference.strategy: continuous
# 开启忽略解析报错
ingestion.ignore-errors: true
# 解析报错发生30次后触发作业失败
ingestion.error-tolerance.max-count: 30
sink:
type: paimon
name: Paimon Sink
catalog.properties.metastore: rest
catalog.properties.uri: dlf_uri
catalog.properties.warehouse: your_warehouse
catalog.properties.token.provider: dlf
#(可选)提交用户名,建议为不同作业设置不同的提交用户以避免冲突
commit.user: your_job_name
#(可选)开启删除向量,提升读取性能
table.properties.deletion-vectors.enabled: true
# 开启脏数据收集器,打印解析失败数据
pipeline:
dirty-data.collector:
name: Logger Dirty Data Collector
type: logger
Baca data JSON
Bagian berikut menjelaskan cara umum untuk membaca data dalam format JSON.
Tentukan ID tabel untuk penguraian
Secara default, ID tabel data JSON adalah nama topik. Anda dapat menentukan nilai field dalam data sebagai ID tabel. Dalam contoh berikut, field db dan tbl ditentukan sebagai ID tabel.
source:
type: kafka
topic: customers
properties.bootstrap.servers: localhost:9092
properties.group.id: ${kafka.group.id}
value.format: json
# (可选)动态识别每条数据Schema,并比对生成Schema变更
schema.inference.strategy: continuous
# 指定字段 db 和 table 作为 table id
value.json.decode.parser-table-id.fields: db,tbl
sink:
type: paimon
name: Paimon Sink
catalog.properties.metastore: rest
catalog.properties.uri: dlf_uri
catalog.properties.warehouse: your_warehouse
catalog.properties.token.provider: dlf
#(可选)提交用户名,建议为不同作业设置不同的提交用户以避免冲突
commit.user: your_job_name
#(可选)开启删除向量,提升读取性能
table.properties.deletion-vectors.enabled: true
# 开启脏数据收集器,打印解析失败数据
pipeline:
dirty-data.collector:
name: Logger Dirty Data Collector
type: logger
Tentukan tipe field
Tipe field diinfer dari nilai field. Beberapa tipe yang diinfer mungkin bukan tipe yang Anda harapkan. Anda dapat menentukan tipe tetap untuk field tertentu dan melewati inferensi dan evolusi tipe field tersebut selanjutnya.
Konfigurasi berikut memperbaiki tipe empat field:
source:
type: kafka
topic: customers
properties.bootstrap.servers: localhost:9092
properties.group.id: ${kafka.group.id}
value.format: json
# 指定字段 db 和 tbl 作为 table id
value.json.decode.parser-table-id.fields: db,tbl
# 固定指定部分字段类型
value.json.infer-schema.fixed-types: 'db STRING, tbl STRING, id BIGINT, amount DECIMAL(18, 2)'
# 允许未声明字段继续动态推断
schema.inference.strategy: continuous
sink:
type: paimon
name: Paimon Sink
catalog.properties.metastore: rest
catalog.properties.uri: dlf_uri
catalog.properties.warehouse: your_warehouse
catalog.properties.token.provider: dlf
#(可选)提交用户名,建议为不同作业设置不同的提交用户以避免冲突
commit.user: your_job_name
#(可选)开启删除向量,提升读取性能
table.properties.deletion-vectors.enabled: true
# 开启脏数据收集器,打印解析失败数据
pipeline:
dirty-data.collector:
name: Logger Dirty Data Collector
type: logger
Baca data Canal JSON
Bagian berikut menjelaskan cara umum untuk membaca data dalam format Canal JSON.
Kebijakan inferensi tipe
Secara default, saat data Canal JSON dibaca, nilai data digunakan untuk menginfer tipe Schema.
source:
type: kafka
topic: customers
properties.bootstrap.servers: localhost:9092
properties.group.id: ${kafka.group.id}
value.format: canal-json
# (可选)动态识别每条数据Schema,并比对生成Schema变更
schema.inference.strategy: continuous
sink:
type: paimon
name: Paimon Sink
catalog.properties.metastore: rest
catalog.properties.uri: dlf_uri
catalog.properties.warehouse: your_warehouse
catalog.properties.token.provider: dlf
#(可选)提交用户名,建议为不同作业设置不同的提交用户以避免冲突
commit.user: your_job_name
#(可选)开启删除向量,提升读取性能
table.properties.deletion-vectors.enabled: true
# 开启脏数据收集器,打印解析失败数据
pipeline:
dirty-data.collector:
name: Logger Dirty Data Collector
type: logger
Anda juga dapat menginfer tipe berdasarkan informasi Schema (tipe sql atau tipe mysql) yang dicatat dalam data Canal JSON. Contoh berikut menggunakan tipe mysql untuk menginfer Schema.
source:
type: kafka
topic: customers
properties.bootstrap.servers: localhost:9092
properties.group.id: ${kafka.group.id}
value.format: canal-json
# (可选)动态识别每条数据Schema,并比对生成Schema变更
schema.inference.strategy: continuous
# 使用 mysql type 信息推导 Schema,也可以配置为 SQL_TYPE 通过 sql type 推导
value.canal-json.infer-schema.strategy: MYSQL_TYPE
sink:
type: paimon
name: Paimon Sink
catalog.properties.metastore: rest
catalog.properties.uri: dlf_uri
catalog.properties.warehouse: your_warehouse
catalog.properties.token.provider: dlf
#(可选)提交用户名,建议为不同作业设置不同的提交用户以避免冲突
commit.user: your_job_name
#(可选)开启删除向量,提升读取性能
table.properties.deletion-vectors.enabled: true
# 开启脏数据收集器,打印解析失败数据
pipeline:
dirty-data.collector:
name: Logger Dirty Data Collector
type: logger
Baca data Debezium JSON
Contoh berikut membaca topik customers dan menulis data ke Data Lake Formation:
source:
type: kafka
topic: customers
properties.bootstrap.servers: localhost:9092
properties.group.id: ${kafka.group.id}
value.format: debezium-json
# (可选)动态识别每条数据Schema,并比对生成Schema变更
schema.inference.strategy: continuous
sink:
type: paimon
name: Paimon Sink
catalog.properties.metastore: rest
catalog.properties.uri: dlf_uri
catalog.properties.warehouse: your_warehouse
catalog.properties.token.provider: dlf
#(可选)提交用户名,建议为不同作业设置不同的提交用户以避免冲突
commit.user: your_job_name
#(可选)开启删除向量,提升读取性能
table.properties.deletion-vectors.enabled: true
Sinkronisasi log biner MySQL mentah ke Kafka
Ingesti data Flink CDC mendukung sinkronisasi konten log biner MySQL mentah ke Canal JSON. Pekerjaan berikut menyinkronkan log biner beberapa tabel ke topik order_dw_tables.
source:
type: mysql
hostname: #{hostname}
port: 3306
username: #{username}
password: #{password}
tables: order_dw.\.*
server-id: 28601-28604
#(可选)同步增量阶段新创建的表的数据
scan.binlog.newly-added-table.enabled: true
#(可选)同步表注释和字段注释
include-comments.enabled: true
#(可选)优先分发无界的分片以避免可能出现的TaskManager OutOfMemory问题
scan.incremental.snapshot.unbounded-chunk-first.enabled: true
#(可选)开启解析过滤,加速读取
scan.only.deserialize.captured.tables.changelog.enabled: true
# 在 Canal JSON 中补充 mysqlType、sqlType、sql、isDdl 等信息
include-binlog-meta.enable: true
sink:
type: kafka
properties.bootstrap.servers: localhost:9092
topic: order_dw_tables
# Kafka value 使用 Canal JSON changelog 格式
value.format: canal-json
# 指定序列化日期类型数据时使用的格式
value.canal-json.timestamp-format.standard: SQL
# 数据统一写入分区0,保证binlog顺序
partition.strategy: all-to-zero
Fitur lanjutan
Kebijakan penguraian dan sinkronisasi perubahan Schema
Konektor Kafka memelihara Schema semua tabel yang saat ini diketahui.
Inisialisasi informasi Schema tabel
Informasi Schema tabel mencakup informasi field dan tipe data, informasi database dan tabel, serta informasi kunci primer. Ketiga jenis informasi ini diinisialisasi sebagai berikut:
-
Informasi field dan tipe data
Pekerjaan ingesti data dapat secara otomatis menginfer field dan tipe data tabel dari data. Namun, dalam beberapa skenario, Anda mungkin ingin menentukan field dan tipe tabel tertentu. Berdasarkan granularitas penentuan tipe field, informasi Schema tabel dapat diinisialisasi dengan menggunakan salah satu dari tiga kebijakan berikut:
-
Sepenuhnya diinferensikan oleh program
Sebelum data Kafka dibaca, konektor Kafka mencoba mengonsumsi hingga jumlah pesan yang ditentukan oleh scan.max.pre.fetch.records dari setiap partisi terlebih dahulu, mengurai Schema setiap record, lalu menggabungkan Schema tersebut untuk menginisialisasi informasi Schema tabel. Sebelum data benar-benar dikonsumsi, event pembuatan tabel yang sesuai dihasilkan berdasarkan Schema yang diinisialisasi.
Untuk format Debezium JSON dan Canal JSON, informasi tabel terdapat dalam pesan individual. Pesan scan.max.pre.fetch.records yang dikonsumsi terlebih dahulu mungkin berisi data beberapa tabel, sehingga jumlah record yang dikonsumsi terlebih dahulu tidak dapat ditentukan untuk setiap tabel. Pra-konsumsi dan inisialisasi Schema hanya dilakukan sekali sebelum pesan setiap partisi benar-benar dikonsumsi dan diproses. Jika data tabel baru tiba nanti, Schema tabel yang diurai dari record pertama tabel tersebut digunakan sebagai Schema awal, dan pra-konsumsi serta inisialisasi tidak dilakukan lagi untuk tabel tersebut.
Data satu tabel yang didistribusikan di beberapa partisi hanya didukung pada VVR 8.0.11 dan yang lebih baru. Dalam skenario ini, Anda harus mengatur debezium-json.distributed-tables atau canal-json.distributed-tables ke true.
-
Tentukan Schema tabel awal
Dalam beberapa skenario, Anda mungkin ingin menentukan Schema tabel awal sendiri, misalnya saat Anda menulis data Kafka ke tabel downstream yang telah dibuat sebelumnya. Dalam kasus ini, Anda dapat menambahkan parameter scan.value.initial-schemas.ddls untuk menentukan Schema tabel awal. Contoh berikut menunjukkan konfigurasi:
source:
type: kafka
name: Kafka Source
properties.bootstrap.servers: host:9092
topic: test-topic
value.format: json
scan.startup.mode: earliest-offset
# 使用数据中的 db、table 字段作为 Table ID
json.decode.parser-table-id.fields: db,table
# 设置初始表结构
scan.value.initial-schemas.ddls: CREATE TABLE db1.t1 (id BIGINT, name VARCHAR(10)); CREATE TABLE db1.t2 (id BIGINT);
Pernyataan CREATE TABLE harus konsisten dengan Schema tabel tujuan. Dalam contoh ini, tipe awal field id dalam tabel db1.t1 ditentukan sebagai BIGINT dan tipe awal field name ditentukan sebagai VARCHAR(10), serta tipe awal field id dalam tabel db1.t2 ditentukan sebagai BIGINT.
Pernyataan CREATE TABLE menggunakan sintaksis Flink SQL.
-
Tentukan tipe tetap untuk field
Dalam beberapa skenario, Anda mungkin ingin memperbaiki tipe data field tertentu. Misalnya, Anda mungkin ingin field tertentu yang dapat diinfer sebagai tipe TIMESTAMP dikirimkan sebagai string. Dalam kasus ini, Anda dapat menambahkan parameter json.infer-schema.fixed-types untuk menentukan Schema tabel awal. Ini hanya berlaku saat format pesan adalah json. Contoh berikut menunjukkan konfigurasi:
source:
type: kafka
name: Kafka Source
properties.bootstrap.servers: host:9092
topic: test-topic
value.format: json
scan.startup.mode: earliest-offset
# 设置特定字段始终为固定类型
json.infer-schema.fixed-types: id BIGINT, name VARCHAR(10)
scan.max.pre.fetch.records: 0
Dalam contoh ini, tipe semua field id diperbaiki sebagai BIGINT dan tipe semua field name diperbaiki sebagai VARCHAR(10).
Tipe yang digunakan di sini sama dengan tipe data Flink SQL.
-
Informasi database dan tabel
-
Untuk format Canal JSON dan Debezium JSON, informasi tabel diurai dari pesan individual dan mencakup nama database dan nama tabel.
-
Untuk format JSON, informasi tabel hanya berisi nama tabel secara default, yaitu nama topik tempat data berada. Jika data Anda berisi informasi database dan tabel, Anda dapat menggunakan json.infer-schema.fixed-types untuk menentukan field yang berisi informasi database dan tabel. Field ini dipetakan ke nama database dan nama tabel. Contoh berikut menunjukkan konfigurasi:
source: type: kafka name: Kafka Source properties.bootstrap.servers: host:9092 topic: test-topic value.format: json scan.startup.mode: earliest-offset # 使用 col1 字段中的值作为库名,使用 col2 字段中的值作为表名 json.decode.parser-table-id.fields: col1,col2Dalam contoh ini, setiap record dikirimkan ke tabel yang nama databasenya adalah nilai field col1 dan nama tabelnya adalah nilai field col2.
-
-
Informasi kunci primer
-
Untuk format Canal JSON, kunci primer tabel didefinisikan berdasarkan field pkNames dalam JSON.
-
Untuk format Debezium JSON dan JSON, JSON tidak berisi informasi kunci primer. Anda dapat menggunakan aturan transform untuk menambahkan kunci primer ke tabel secara manual:
transform: - source-table: \.*.\.* projection: \* primary-keys: key1, key2
-
Penguraian Schema dan perubahan Schema
Setelah Schema tabel diinisialisasi, jika schema.inference.strategy diatur ke static, konektor Kafka mengurai nilai setiap pesan berdasarkan Schema tabel awal dan tidak menghasilkan event perubahan Schema. Jika schema.inference.strategy diatur ke continuous, konektor Kafka mengurai nilai setiap pesan Kafka, memperoleh kolom fisik pesan, dan membandingkannya dengan Schema yang saat ini dipelihara. Jika Schema yang diurai tidak konsisten dengan Schema saat ini, konektor mencoba menggabungkan Schema dan menghasilkan event perubahan Schema tabel yang sesuai. Aturan penggabungan adalah sebagai berikut:
-
Jika kolom fisik yang diurai berisi field yang tidak ada dalam Schema saat ini, field tersebut ditambahkan ke Schema dan event add-nullable-column dihasilkan.
-
Jika kolom fisik yang diurai tidak berisi field yang sudah ada dalam Schema saat ini, field tersebut dipertahankan dan data kolom tersebut diisi dengan NULL. Tidak ada event drop-column yang dihasilkan.
-
Jika keduanya berisi kolom dengan nama yang sama, kolom tersebut ditangani sebagai berikut:
-
Saat tipenya sama tetapi presisinya berbeda, tipe dengan presisi lebih tinggi digunakan dan event column-type-change dihasilkan.
-
Saat tipenya berbeda, node induk paling rendah dalam struktur pohon berikut digunakan sebagai tipe kolom, dan event column-type-change dihasilkan.

-
-
Kebijakan perubahan Schema berikut saat ini didukung:
-
Tambah kolom: Kolom yang sesuai ditambahkan ke akhir Schema saat ini, data kolom baru disinkronkan, dan kolom baru diatur sebagai kolom nullable.
-
Hapus kolom: Tidak ada event drop-column yang dihasilkan. Sebagai gantinya, data kolom secara otomatis diisi dengan nilai NULL.
-
Ganti nama kolom: Ini dianggap sebagai menambah kolom dan menghapus kolom. Kolom yang diganti namanya ditambahkan ke akhir Schema saat ini, dan data kolom sebelum penggantian nama diisi dengan nilai NULL.
-
Ubah tipe kolom:
-
Untuk sistem downstream yang mendukung perubahan tipe kolom, pekerjaan ingesti data mendukung perubahan tipe kolom biasa setelah sink downstream mendukung penanganan perubahan tipe kolom. Misalnya, kolom dapat diubah dari tipe INT ke tipe BIGINT. Perubahan semacam ini bergantung pada aturan perubahan tipe kolom yang didukung oleh sink downstream. Tabel sink yang berbeda mendukung aturan yang berbeda. Rujuk dokumentasi tabel sink untuk aturan perubahan tipe kolom yang didukungnya.
-
Untuk sistem downstream yang tidak mendukung perubahan tipe kolom, seperti Hologres, Anda dapat menggunakan Pelebaran tipe. Dalam kasus ini, tabel dengan tipe yang lebih luas dibuat di sistem downstream saat pekerjaan dimulai. Saat terjadi perubahan tipe kolom, pekerjaan memeriksa apakah sink downstream dapat menerima perubahan tersebut, yang mengimplementasikan dukungan toleran untuk perubahan tipe kolom.
-
-
-
Perubahan Schema berikut tidak didukung:
-
Perubahan pada kendala seperti kunci primer atau indeks.
-
Perubahan dari NOT NULL ke NULLABLE.
-
-
Penguraian schema Canal JSON
Data Canal JSON mungkin berisi field sqlType opsional yang mencatat informasi tipe yang tepat dari kolom data. Untuk memperoleh Schema yang lebih akurat, Anda dapat mengatur canal-json.infer-schema.strategy ke SQL_TYPE untuk menggunakan tipe dalam sqlType. Pemetaan tipe adalah sebagai berikut:
Tipe JDBC
Kode Jenis
Tipe CDC
BIT
-7
BOOLEAN
BOOLEAN
16
TINYINT
-6
TINYINT
SMALLINT
-5
SMALLINT
INTEGER
4
INT
BIGINT
-5
BIGINT
DECIMAL
3
DECIMAL(38,18)
NUMERIC
2
REAL
7
FLOAT
FLOAT
6
DOUBLE
8
DOUBLE
BINARY
-2
BYTES
VARBINARY
-3
LONGVARBINARY
-4
BLOB
2004
DATE
91
DATE
TIME
92
TIME
TIMESTAMP
93
TIMESTAMP
CHAR
1
STRING
VARCHAR
12
LONGVARCHAR
-1
Tipe lainnya
Toleransi dan pengumpulan data kotor
Dalam beberapa kasus, sumber data Kafka Anda mungkin berisi data rusak (data kotor). Untuk mencegah pekerjaan sinkronisasi gagal dan sering dimulai ulang karena data kotor ini, Anda dapat mengonfigurasi pekerjaan untuk mengabaikan data abnormal tersebut. Contoh berikut menunjukkan konfigurasi:
source:
type: kafka
name: Kafka Source
properties.bootstrap.servers: host:9092
topic: test-topic
value.format: json
scan.startup.mode: earliest-offset
# 开启脏数据容忍功能
ingestion.ignore-errors: true
# 容忍 1000 条脏数据
ingestion.error-tolerance.max-count: 1000
Konfigurasi ini mengabaikan hingga 1.000 record kotor, yang memungkinkan pekerjaan berjalan normal saat terdapat sedikit data kotor. Saat jumlah record kotor melebihi ambang batas ini, pekerjaan gagal, yang mengingatkan Anda untuk memvalidasi data.
Jika Anda tidak ingin pekerjaan gagal karena data kotor, gunakan konfigurasi berikut:
source:
type: kafka
name: Kafka Source
properties.bootstrap.servers: host:9092
topic: test-topic
value.format: json
scan.startup.mode: earliest-offset
# 开启脏数据容忍功能
ingestion.ignore-errors: true
# 容忍所有的脏数据
ingestion.error-tolerance.max-count: -1
Kebijakan toleransi data kotor mencegah pekerjaan gagal sering karena data abnormal. Anda mungkin juga ingin mempelajari lebih lanjut tentang data kotor untuk menyesuaikan perilaku produsen data Kafka. Dengan mengikuti proses yang dijelaskan dalam Pengumpulan data kotor, Anda dapat melihat data kotor pekerjaan di log TaskManager. Contoh berikut menunjukkan konfigurasi:
source:
type: kafka
name: Kafka Source
properties.bootstrap.servers: host:9092
topic: test-topic
value.format: json
scan.startup.mode: earliest-offset
# 开启脏数据容忍功能
ingestion.ignore-errors: true
# 容忍所有的脏数据
ingestion.error-tolerance.max-count: -1
pipeline:
dirty-data.collector:
# 将脏数据写入 TaskManager 的日志文件中
type: logger
Kebijakan pemetaan antara nama tabel dan topik
Saat Kafka digunakan sebagai sink pekerjaan ingesti data, format pesan Kafka (debezium-json atau canal-json) juga berisi informasi nama tabel. Saat pesan Kafka dikonsumsi nanti, nama tabel dalam data biasanya digunakan sebagai nama tabel aktual alih-alih nama topik. Oleh karena itu, Anda harus hati-hati mengonfigurasi kebijakan pemetaan antara nama tabel dan topik.
Asumsikan bahwa dua tabel mydb.mytable1 dan mydb.mytable2 di MySQL perlu disinkronkan. Kebijakan konfigurasi berikut tersedia:
1. Jangan mengonfigurasi kebijakan pemetaan apa pun
Tanpa kebijakan pemetaan apa pun, setiap tabel ditulis ke topik yang sesuai yang dinamai dalam format nama database.nama tabel. Oleh karena itu, data mydb.mytable1 ditulis ke topik bernama mydb.mytable1, dan data mydb.mytable2 ditulis ke topik bernama mydb.mytable2. Contoh berikut menunjukkan konfigurasi:
source:
type: mysql
name: MySQL Source
hostname: ${secret_values.mysql.hostname}
port: ${mysql.port}
username: ${secret_values.mysql.username}
password: ${secret_values.mysql.password}
tables: mydb.mytable1,mydb.mytable2
server-id: 8601-8604
sink:
type: kafka
name: Kafka Sink
properties.bootstrap.servers: ${kafka.bootstraps.server}
2. Konfigurasikan aturan route untuk pemetaan (tidak direkomendasikan)
Dalam banyak skenario, Anda tidak ingin topik tujuan berada dalam format nama database.nama tabel dan ingin menulis data ke topik tertentu. Dalam kasus ini, Anda dapat mengonfigurasi aturan route untuk pemetaan. Contoh berikut menunjukkan konfigurasi:
source:
type: mysql
name: MySQL Source
hostname: ${secret_values.mysql.hostname}
port: ${mysql.port}
username: ${secret_values.mysql.username}
password: ${secret_values.mysql.password}
tables: mydb.mytable1,mydb.mytable2
server-id: 8601-8604
sink:
type: kafka
name: Kafka Sink
properties.bootstrap.servers: ${kafka.bootstraps.server}
route:
- source-table: mydb.mytable1,mydb.mytable2
sink-table: mytable1
Dalam kasus ini, semua data dari mydb.mytable1 dan mydb.mytable2 ditulis ke satu topik mytable1.
Namun, saat aturan route digunakan untuk mengubah nama topik tujuan, informasi nama tabel dalam pesan Kafka (dalam format debezium-json atau canal-json) juga berubah. Dalam kasus ini, semua nama tabel dalam pesan Kafka menjadi mytable1, yang dapat menyebabkan hasil yang tidak diinginkan saat sistem lain mengonsumsi pesan Kafka dari topik ini.
3. Konfigurasikan parameter sink.tableId-to-topic.mapping untuk pemetaan (direkomendasikan)
Untuk mengonfigurasi aturan pemetaan antara nama tabel dan topik sambil mempertahankan informasi nama tabel sumber, Anda dapat menggunakan parameter sink.tableId-to-topic.mapping. Contoh berikut menunjukkan konfigurasi:
source:
type: mysql
name: MySQL Source
hostname: ${secret_values.mysql.hostname}
port: ${mysql.port}
username: ${secret_values.mysql.username}
password: ${secret_values.mysql.password}
tables: mydb.mytable1,mydb.mytable2
server-id: 8601-8604
sink:
type: kafka
name: Kafka Sink
properties.bootstrap.servers: ${kafka.bootstraps.server}
sink.tableId-to-topic.mapping: mydb.mytable1,mydb.mytable2:mytable
Atau
source:
type: mysql
name: MySQL Source
hostname: ${secret_values.mysql.hostname}
port: ${mysql.port}
username: ${secret_values.mysql.username}
password: ${secret_values.mysql.password}
tables: mydb.mytable1,mydb.mytable2
server-id: 8601-8604
sink:
type: kafka
name: Kafka Sink
properties.bootstrap.servers: ${kafka.bootstraps.server}
sink.tableId-to-topic.mapping: mydb.mytable1:mytable;mydb.mytable2:mytable
Dalam kasus ini, semua data dari mydb.mytable1 dan mydb.mytable2 ditulis ke satu topik mytable1, dan informasi nama tabel dalam pesan Kafka (dalam format debezium-json atau canal-json) tetap mydb.mytable1 atau mydb.mytable2. Saat sistem lain mengonsumsi pesan Kafka dari topik ini, mereka dapat memperoleh informasi nama tabel sumber dengan benar.
Implementasi Konverter JSON
Data JSON mungkin tidak selalu memiliki format yang sepenuhnya konsisten, atau beberapa persyaratan pemrosesan tidak dapat dipenuhi oleh modul Transform. Misalnya, Anda mungkin perlu menggabungkan dua kolom menjadi kolom baru, menghapus dua kolom asli, dan memastikan bahwa sinkronisasi perubahan Schema masih berfungsi. Untuk memproses data JSON secara lebih fleksibel, Anda dapat mengimplementasikan antarmuka KafkaPayloadConverter untuk memodifikasi data JSON sebelum framework memprosesnya. Untuk menggunakan fitur ini, tambahkan konfigurasi json.decode.converter-class dan atur ke nama lengkap kelas implementasi.
Proses implementasi dan penggunaan Konverter JSON adalah sebagai berikut:
-
Implementasikan antarmuka KafkaPayloadConverter dan paketkan. Repositori demo open source sudah menyediakan beberapa contoh implementasi.
-
Tambahkan file yang telah dipaketkan ke dependensi tambahan pekerjaan ingesti data.
-
Di modul Sumber, tentukan konfigurasi json.decode.converter-class. Contoh berikut menggunakan kelas ArrayElementExtractorConverter dari proyek open source:
source: type: kafka properties.bootstrap.servers: localhost:9092 topic: my_cdc_topic properties.group.id: flink-cdc-group scan.startup.mode: earliest-offset value.format: json # 对 Kafka 消息 value 的 JSON 数据应用自定义转换器 value.json.decode.converter-class: org.apache.flink.cdc.connectors.kafka.ArrayElementExtractorConverter -
Deploy dan jalankan pekerjaan.
Format kunci dan format nilai
Pada VVR 11.5 dan yang lebih awal, konfigurasi format kunci dan format nilai tidak dapat dibedakan. Saat format yang sama digunakan, konfigurasi format berlaku untuk kedua format kunci dan nilai.
VVR 11.6 dan yang lebih baru mengoptimalkan perilaku ini. Mengambil format json sebagai contoh, item konfigurasi format diteruskan sebagai berikut:
-
Awalan format (seperti json.infer-schema.primitive-as-string): Untuk menjaga kompatibilitas versi, konfigurasi dengan awalan format berlaku untuk kedua format kunci dan nilai secara default.
-
Awalan kunci ditambah awalan format (seperti key.json.infer-schema.primitive-as-string): Konfigurasi dengan awalan kunci ditambah awalan format hanya berlaku untuk format kunci dan memiliki prioritas lebih tinggi daripada konfigurasi dengan awalan format.
-
Awalan nilai ditambah awalan format (seperti value.json.infer-schema.primitive-as-string): Konfigurasi dengan awalan nilai ditambah awalan format hanya berlaku untuk format nilai dan memiliki prioritas lebih tinggi daripada konfigurasi dengan awalan format.
Dalam konfigurasi sumber Kafka berikut, Konverter JSON berbeda dikonfigurasi untuk format kunci dan format nilai.
source:
type: kafka
properties.bootstrap.servers: localhost:9092
topic: my_cdc_topic
properties.group.id: flink-cdc-group
scan.startup.mode: earliest-offset
key.format: json
value.format: json
# 对 Kafka 消息 key 的 JSON 数据应用自定义转换器
key.json.decode.converter-class: com.example.KeyExampleConverter
# 对 Kafka 消息 value 的 JSON 数据应用自定义转换器
value.json.decode.converter-class: com.example.ValueExampleConverter