Setelah mengonfigurasi Katalog JSON Kafka, pekerjaan Realtime Compute for Apache Flink Anda dapat langsung mengakses topik JSON di kluster Kafka tanpa perlu mendefinisikan skema. Topik ini menjelaskan cara membuat, melihat, dan menghapus Katalog JSON Kafka.
Latar Belakang
Katalog JSON Kafka melakukan inferensi skema topik dengan secara otomatis mengurai pesan berformat JSON. Hal ini memungkinkan Anda mengakses bidang tertentu dari pesan dalam Flink SQL tanpa harus mendeklarasikan skema untuk tabel Kafka. Katalog JSON Kafka memiliki fitur-fitur berikut:
-
Nama tabel dalam katalog sesuai dengan nama topik Kafka, sehingga menghilangkan kebutuhan pendaftaran manual tabel menggunakan pernyataan DDL dan meningkatkan efisiensi serta akurasi pengembangan.
-
Anda dapat langsung menggunakan tabel dari katalog JSON Kafka sebagai tabel sumber dalam pekerjaan Flink SQL.
-
Anda dapat menggunakan katalog JSON Kafka dengan pernyataan CREATE TABLE AS (CTAS) untuk menyinkronkan data selama perubahan skema.
Topik ini menjelaskan cara mengelola katalog JSON Kafka:
Batasan
-
Katalog JSON Kafka hanya mendukung topik dengan pesan dalam format JSON.
-
Hanya mesin komputasi Flink VVR 6.0.2 atau yang lebih baru yang mendukung konfigurasi katalog JSON Kafka.
CatatanUntuk menggunakan katalog JSON Kafka pada VVR 4.x, upgrade pekerjaan Anda ke VVR 6.0.2 atau yang lebih baru.
-
Anda tidak dapat menggunakan pernyataan DDL untuk memodifikasi katalog JSON Kafka yang sudah ada.
-
Anda hanya dapat melakukan kueri terhadap tabel data. Anda tidak dapat membuat, memodifikasi, atau menghapus database atau tabel.
CatatanDalam skenario CREATE DATABASE AS SELECT (CDAS) atau CREATE TABLE AS SELECT (CTAS) yang menggunakan katalog JSON Kafka, topik dapat dibuat secara otomatis.
-
Katalog JSON Kafka tidak dapat membaca dari atau menulis ke kluster Kafka yang menggunakan autentikasi SSL atau SASL.
-
Anda dapat menggunakan tabel dari katalog JSON Kafka secara langsung sebagai tabel sumber dalam pekerjaan Flink SQL, tetapi tidak sebagai tabel sink atau tabel dimensi lookup.
-
Karena ApsaraMQ for Kafka saat ini tidak mendukung penghapusan Group menggunakan antarmuka yang sama seperti Kafka open-source, Anda harus menentukan aliyun.kafka.instanceId, aliyun.kafka.accessKeyId, aliyun.kafka.accessKeySecret, aliyun.kafka.endpoint, dan aliyun.kafka.regionId saat membuat katalog JSON Kafka agar Group ID dapat dihapus secara otomatis. Untuk informasi lebih lanjut, lihat Perbandingan dengan Kafka open-source.
Catatan Penggunaan
Katalog JSON Kafka melakukan inferensi skema tabel dengan mengurai data sampel. Ketika sebuah topik berisi format data yang tidak konsisten, katalog secara default menyimpan semua kolom dan mengembalikan skema terluas. Perubahan format data dalam topik dapat mengubah skema hasil inferensi, yang berpotensi menyebabkan kegagalan pekerjaan saat restart.
Sebagai contoh, ketika me-restart pekerjaan Flink SQL yang membaca dari tabel Katalog JSON Kafka dari titik simpan (savepoint), katalog mungkin melakukan inferensi skema baru yang berbeda dari skema sebelumnya. Rencana eksekusi pekerjaan didasarkan pada skema lama, yang dapat menyebabkan error waktu proses seperti kondisi filter yang tidak sesuai atau akses bidang yang salah. Untuk mencegah masalah ini, gunakan pernyataan CREATE TEMPORARY TABLE dalam pekerjaan Flink SQL Anda untuk mendefinisikan tabel Kafka dengan skema tetap.
Buat katalog JSON Kafka
-
Pada halaman Scripts, masukkan pernyataan untuk membuat katalog JSON Kafka di editor SQL.
-
Untuk kluster Kafka yang dikelola sendiri atau kluster EMR on ECS Kafka
CREATE CATALOG <YourCatalogName> WITH( 'type'='kafka', 'properties.bootstrap.servers'='<brokers>', 'format'='json', 'default-database'='<dbName>', 'key.fields-prefix'='<keyPrefix>', 'value.fields-prefix'='<valuePrefix>', 'timestamp-format.standard'='<timestampFormat>', 'infer-schema.flatten-nested-columns.enable'='<flattenNestedColumns>', 'infer-schema.primitive-as-string'='<primitiveAsString>', 'infer-schema.parse-key-error.field-name'='<parseKeyErrorFieldName>', 'infer-schema.compacted-topic-as-upsert-table'='true', 'max.fetch.records'='100' ); -
Untuk ApsaraMQ for Kafka
CREATE CATALOG <YourCatalogName> WITH( 'type'='kafka', 'properties.bootstrap.servers'='<brokers>', 'format'='json', 'default-database'='<dbName>', 'key.fields-prefix'='<keyPrefix>', 'value.fields-prefix'='<valuePrefix>', 'timestamp-format.standard'='<timestampFormat>', 'infer-schema.flatten-nested-columns.enable'='<flattenNestedColumns>', 'infer-schema.primitive-as-string'='<primitiveAsString>', 'infer-schema.parse-key-error.field-name'='<parseKeyErrorFieldName>', 'infer-schema.compacted-topic-as-upsert-table'='true', 'max.fetch.records'='100', 'aliyun.kafka.accessKeyId'='<aliyunAccessKeyId>', 'aliyun.kafka.accessKeySecret'='<aliyunAccessKeySecret>', 'aliyun.kafka.instanceId'='<aliyunKafkaInstanceId>', 'aliyun.kafka.endpoint'='<aliyunKafkaEndpoint>', 'aliyun.kafka.regionId'='<aliyunKafkaRegionId>' );
Parameter
Type
Deskripsi
Wajib
Keterangan
YourCatalogName
String
Nama katalog JSON Kafka.
Ya
Tentukan nama kustom.
PentingGanti placeholder dengan nama katalog Anda dan hapus tanda kurung sudut (<>) untuk mencegah error sintaksis.
type
String
Jenis katalog.
Ya
Nilainya harus 'kafka'.
properties.bootstrap.servers
String
Alamat broker Kafka.
Ya
Tentukan satu atau beberapa alamat dalam format
host1:port1,host2:port2,host3:port3.Pisahkan beberapa alamat dengan koma (,).
format
String
Format pesan Kafka.
Ya
Saat ini, hanya 'json' yang didukung. Flink mengurai pesan berformat JSON untuk melakukan inferensi skema.
default-database
String
Nama database default untuk katalog.
Tidak
Katalog menggunakan pengenal tiga bagian untuk menemukan tabel: catalog_name.db_name.table_name. Parameter ini mengatur db_name. Karena Kafka tidak memiliki database, Anda dapat menggunakan string apa pun untuk merepresentasikan kluster.
key.fields-prefix
String
Awalan kustom yang ditambahkan ke nama bidang yang diurai dari kunci pesan untuk mencegah konflik penamaan.
Tidak
Nilai default adalah 'key_'. Misalnya, jika bidang dalam kunci bernama 'a', namanya dalam skema tabel menjadi 'key_a'.
CatatanNilai key.fields-prefix tidak boleh menjadi awalan dari nilai value.fields-prefix. Misalnya, jika value.fields-prefix diatur ke 'test1_value_', Anda tidak dapat mengatur key.fields-prefix ke 'test1_'.
value.fields-prefix
String
Awalan kustom yang ditambahkan ke nama bidang yang diurai dari nilai pesan untuk mencegah konflik penamaan.
Tidak
Nilai default adalah 'value_'. Misalnya, jika bidang dalam nilai bernama 'b', namanya dalam skema tabel menjadi 'value_b'.
CatatanNilai value.fields-prefix tidak boleh menjadi awalan dari nilai key.fields-prefix. Misalnya, jika key.fields-prefix diatur ke 'test2_value_', Anda tidak dapat mengatur value.fields-prefix ke 'test2_'.
timestamp-format.standard
String
Format untuk mengurai bidang timestamp dalam pesan JSON. Flink pertama-tama mencoba mengurai timestamp menggunakan format ini. Jika penguraian gagal, sistem secara otomatis mencoba format standar lainnya.
Tidak
Nilai yang valid:
-
SQL (default)
-
ISO-8601
infer-schema.flatten-nested-columns.enable
Boolean
Menentukan apakah akan memperluas kolom bersarang secara rekursif dalam nilai pesan JSON.
Tidak
Nilai yang valid:
-
true: Memperluas kolom bersarang secara rekursif.
Untuk kolom yang diperluas, Flink menggunakan jalur akses sebagai nama kolom. Misalnya, untuk kolom 'col' dalam
{"nested": {"col": true}}, nama kolom yang diperluas adalah 'nested.col'.CatatanJika parameter ini diatur ke true, kami sarankan Anda menggunakannya dengan pernyataan CREATE TABLE AS (CTAS). Pernyataan DML lain tidak mendukung perluasan otomatis kolom bersarang.
-
false (default): Memperlakukan tipe bersarang sebagai STRING.
infer-schema.primitive-as-string
Boolean
Menentukan apakah akan melakukan inferensi semua tipe primitif dalam nilai pesan JSON sebagai STRING.
Tidak
Nilai yang valid:
-
true: Melakukan inferensi semua tipe primitif sebagai STRING.
-
false (default): Melakukan inferensi tipe berdasarkan aturan standar.
infer-schema.parse-key-error.field-name
String
Jika kunci pesan yang tidak kosong tidak dapat diurai sebagai JSON, Flink menambahkan bidang VARBINARY ke skema untuk menyimpan data kunci mentah. Nama bidang baru ini merupakan kombinasi dari key.fields-prefix dan nilai parameter ini.
Tidak
Nilai default adalah 'col'. Misalnya, jika nilai pesan diurai menjadi bidang bernama 'value_name' dan kunci pesan tidak kosong tetapi gagal diurai, skema hasil berisi dua bidang: 'key_col' dan 'value_name'.
infer-schema.compacted-topic-as-upsert-table
Boolean
Menentukan apakah topik diperlakukan sebagai tabel Upsert Kafka. Ini hanya berlaku ketika kebijakan pembersihan topik adalah 'compact' dan kunci pesan tidak kosong.
Tidak
Nilai default adalah true. Parameter ini harus diatur ke true ketika Anda menggunakan pernyataan CTAS atau CDAS untuk menyinkronkan data ke instans ApsaraMQ for Kafka.
CatatanParameter ini hanya didukung di VVR 6.0.2 dan yang lebih baru.
max.fetch.records
Int
Jumlah maksimum pesan yang diambil sebagai sampel untuk inferensi skema.
Tidak
Nilai default adalah 100.
aliyun.kafka.accessKeyId
String
ID AccessKey Akun Alibaba Cloud Anda. Untuk informasi lebih lanjut, lihat Buat Pasangan Kunci Akses.
Tidak
Parameter ini wajib saat Anda menggunakan pernyataan CTAS atau CDAS untuk menyinkronkan data ke instans ApsaraMQ for Kafka.
CatatanParameter ini hanya didukung di VVR 6.0.2 dan yang lebih baru.
aliyun.kafka.accessKeySecret
String
Rahasia AccessKey Akun Alibaba Cloud Anda. Untuk informasi lebih lanjut, lihat Buat Pasangan Kunci Akses.
Tidak
Parameter ini wajib saat Anda menggunakan pernyataan CTAS atau CDAS untuk menyinkronkan data ke instans ApsaraMQ for Kafka.
CatatanParameter ini hanya didukung di VVR 6.0.2 dan yang lebih baru.
aliyun.kafka.instanceId
String
ID instans ApsaraMQ for Kafka. Anda dapat melihat ID instans di halaman detail instans di Konsol ApsaraMQ for Kafka.
Tidak
Parameter ini wajib saat Anda menggunakan pernyataan CTAS atau CDAS untuk menyinkronkan data ke instans ApsaraMQ for Kafka.
CatatanParameter ini hanya didukung di VVR 6.0.2 dan yang lebih baru.
aliyun.kafka.endpoint
String
Titik akhir API untuk ApsaraMQ for Kafka. Untuk informasi lebih lanjut, lihat Endpoints.
Tidak
Parameter ini wajib saat Anda menggunakan pernyataan CTAS atau CDAS untuk menyinkronkan data ke instans ApsaraMQ for Kafka.
CatatanParameter ini hanya didukung di VVR 6.0.2 dan yang lebih baru.
aliyun.kafka.regionId
String
ID Wilayah instans tempat topik berada. Untuk informasi lebih lanjut, lihat Endpoints.
Tidak
Parameter ini wajib saat Anda menggunakan pernyataan CTAS atau CDAS untuk menyinkronkan data ke instans ApsaraMQ for Kafka.
CatatanParameter ini hanya didukung di VVR 6.0.2 dan yang lebih baru.
-
-
Pilih pernyataan CREATE CATALOG, lalu klik Run di gutter sebelah kiri.
CREATE CATALOG KafkaCatalog WITH( 'type'='kafka', 'properties.bootstrap.servers'='<brokers>', 'format'='json', 'default-database'='<dbName>', 'key.fields-prefix'='<keyPrefix>', 'value.fields-prefix'='<valuePrefix>', 'timestamp-format.standard'='<timestampFormat>', 'infer-schema.flatten-nested-columns.enable'='<flattenNestedColumns>', 'infer-schema.primitive-as-string'='<primitiveAsString>', 'infer-schema.parse-key-error.field-name'='<parseKeyErrorFieldName>', 'infer-schema.compacted-topic-as-upsert-table'='true', 'max.fetch.records'='100', 'aliyun.kafka.accessKeyId'='<aliyunAccessKeyId>', 'aliyun.kafka.accessKeySecret'='<aliyunAccessKeySecret>', 'aliyun.kafka.instanceId'='<aliyunKafkaInstanceId>', 'aliyun.kafka.endpoint'='<aliyunKafkaEndpoint>', 'aliyun.kafka.regionId'='<aliyunKafkaRegionId>' ); -
Verifikasi bahwa katalog muncul di area Catalogs di sebelah kiri.
Katalog JSON Kafka
-
Masukkan perintah berikut di editor Scripts.
DESCRIBE `${catalog_name}`.`${db_name}`.`${topic_name}`;Parameter
Deskripsi
${catalog_name}
Nama katalog JSON Kafka.
${db_name}
Nama kluster Kafka.
${topic_name}
Nama topik Kafka.
-
Pilih pernyataan DESCRIBE, lalu klik Run di gutter sebelah kiri.
Setelah pernyataan berhasil dijalankan, panel hasil menampilkan definisi bidang untuk tabel sumber Kafka. Output mencakup bidang bisnis
value_from(STRING),value_name(STRING), danvalue_order(BIGINT), serta bidang virtual metadatapartition(INT, primary key),offset(BIGINT, primary key), dantimestamp(TIMESTAMP_LTZ(3)). Untuk bidang virtual metadata ini, kolomextrasmenampilkanMETADATA VIRTUAL.
Gunakan katalog JSON Kafka
-
Sebagai tabel sumber untuk membaca data dari topik Kafka.
INSERT INTO ${other_sink_table} SELECT... FROM `${kafka_catalog}`.`${db_name}`.`${topic_name}`/*+OPTIONS('scan.startup.mode'='earliest-offset')*/;CatatanUntuk menentukan opsi WITH tambahan untuk tabel dalam katalog JSON Kafka, gunakan Petunjuk SQL. Sebagai contoh, pernyataan SQL di atas menggunakan petunjuk untuk mulai mengonsumsi data dari offset paling awal. Untuk informasi lebih lanjut tentang parameter lainnya, lihat Parameter tabel sumber ApsaraMQ for Kafka dan Parameter tabel sink ApsaraMQ for Kafka.
-
Sebagai tabel sumber dengan pernyataan CREATE TABLE AS (CTAS) untuk menyinkronkan data dari topik Kafka ke tabel tujuan.
-
Menyinkronkan satu tabel secara real-time.
CREATE TABLE IF NOT EXISTS `${target_table_name}` WITH(...) AS TABLE `${kafka_catalog}`.`${db_name}`.`${topic_name}` /*+OPTIONS('scan.startup.mode'='earliest-offset')*/; -
Menyinkronkan beberapa tabel dalam satu pekerjaan.
BEGIN STATEMENT SET; CREATE TABLE IF NOT EXISTS `some_catalog`.`some_database`.`some_table0` AS TABLE `kafka-catalog`.`kafka`.`topic0` /*+ OPTIONS('scan.startup.mode'='earliest-offset') */; CREATE TABLE IF NOT EXISTS `some_catalog`.`some_database`.`some_table1` AS TABLE `kafka-catalog`.`kafka`.`topic1` /*+ OPTIONS('scan.startup.mode'='earliest-offset') */; CREATE TABLE IF NOT EXISTS `some_catalog`.`some_database`.`some_table2` AS TABLE `kafka-catalog`.`kafka`.`topic2` /*+ OPTIONS('scan.startup.mode'='earliest-offset') */; END;Dengan katalog JSON Kafka, Anda dapat menyinkronkan beberapa tabel Kafka dalam pekerjaan yang sama. Ketentuan berikut berlaku:
-
Tidak ada tabel Kafka yang dikonfigurasi dengan parameter topic-pattern.
-
Konfigurasi Kafka untuk setiap tabel harus identik, termasuk semua pengaturan properties.* seperti properties.bootstrap.servers dan properties.group.id.
-
Pengaturan scan.startup.mode untuk setiap tabel harus identik dan hanya dapat diatur ke group-offsets, latest-offset, atau earliest-offset.
Sebagai contoh, pada gambar berikut, dua tabel atas memenuhi ketentuan, sedangkan dua tabel bawah melanggarnya.

-
-
Untuk contoh end-to-end penggunaan katalog JSON Kafka, lihat warehousing log waktu nyata.
Hapus katalog JSON Kafka
Menghapus katalog JSON Kafka tidak memengaruhi pekerjaan yang sedang berjalan. Namun, pekerjaan yang menggunakan tabel dari katalog tersebut akan gagal dengan error "table not found" jika Anda deploy ulang atau restart. Lakukan dengan hati-hati.
-
Masukkan perintah berikut di editor SQL pada halaman Scripts.
DROP CATALOG ${catalog_name};Ganti ${catalog_name} dengan nama katalog JSON Kafka Anda.
-
Pilih perintah DROP CATALOG, klik kanan, lalu pilih Run.
-
Verifikasi bahwa katalog tidak lagi tercantum di area Catalogs di sebelah kiri.
Tabel Katalog JSON Kafka
Untuk menyederhanakan penggunaan tabel dari Katalog JSON Kafka, katalog secara otomatis menambahkan parameter konfigurasi default, metadata, dan informasi primary key ke tabel hasil inferensi.
-
Inferensi skema untuk tabel Kafka
Saat melakukan inferensi skema untuk topik, Katalog JSON Kafka mengambil sampel hingga max.fetch.records pesan. Sistem mengurai skema setiap pesan lalu menggabungkan skema-skema tersebut menjadi skema akhir. Aturan penguraian sama dengan aturan dasar yang digunakan ketika Kafka menjadi sumber data untuk pernyataan CREATE TABLE AS (CTAS).
Penting-
Selama inferensi skema, Katalog JSON Kafka membuat kelompok konsumen untuk membaca data dari topik. Katalog memberikan awalan pada nama kelompok konsumen untuk menunjukkan bahwa kelompok tersebut dibuat oleh katalog.
-
Untuk ApsaraMQ for Kafka, kami sarankan menggunakan Katalog JSON Kafka dengan versi 6.0.7 atau yang lebih baru. Versi sebelum 6.0.7 tidak menghapus kelompok konsumen secara otomatis, yang dapat menyebabkan peringatan tentang backlog pesan.
-
Kolom fisik hasil inferensi
Katalog JSON Kafka melakukan inferensi kolom fisik dari kunci pesan dan nilai pesan.
Jika kunci pesan tidak kosong tetapi tidak dapat diurai, katalog mengembalikan kolom VARBINARY. Nama kolom ini merupakan gabungan dari nilai parameter key.fields-prefix dan nilai parameter infer-schema.parse-key-error.field-name.
Setelah mengambil batch pesan Kafka, katalog mengurai setiap pesan dan menggabungkan kolom fisik yang dihasilkan menjadi skema terpadu untuk topik tersebut.
-
Jika kolom fisik hasil inferensi berisi bidang yang tidak ada dalam skema hasil, katalog secara otomatis menambahkannya ke skema hasil.
-
-
Jika tipe sama tetapi presisinya berbeda, sistem menggunakan tipe dengan presisi lebih tinggi.
-
Jika tipenya berbeda, sistem mencari induk umum terendah dalam pohon tipe (ditunjukkan pada gambar di bawah) dan menggunakannya sebagai tipe kolom. Namun, untuk mempertahankan presisi, ketika tipe Decimal dan Float digabung, tipe hasilnya adalah Double.

-
Sebagai contoh, untuk topik Kafka yang berisi tiga pesan berikut, Katalog JSON Kafka melakukan inferensi skema seperti yang ditunjukkan pada gambar di bawah.

-
-
Kolom metadata default
Secara default, Katalog JSON Kafka menambahkan tiga kolom metadata: partition, offset, dan timestamp.
Nama metadata
Nama kolom
Tipe
Deskripsi
partition
partition
INT NOT NULL
ID partisi.
offset
offset
BIGINT NOT NULL
Offset pesan.
timestamp
timestamp
TIMESTAMP_LTZ(3) NOT NULL
Timestamp pesan.
-
Batasan primary key default
Ketika tabel dari Katalog JSON Kafka digunakan sebagai sumber, kolom metadata partition dan offset berfungsi sebagai primary key default untuk memastikan keunikan data.
CatatanJika skema yang diinferensi oleh Katalog JSON Kafka tidak memenuhi kebutuhan Anda, Anda dapat mendeklarasikan tabel temporary dengan skema yang diinginkan menggunakan sintaks CREATE TEMPORARY TABLE ... LIKE. Sebagai contoh, jika data JSON Anda berisi bidang
tsdengan format '2023-01-01 12:00:01', Katalog JSON Kafka secara otomatis melakukan inferensi bidangtssebagai tipe TIMESTAMP. Untuk menggunakan bidangtssebagai tipe STRING, deklarasikan tabel sebagai berikut. Karena konfigurasi default menambahkan awalanvalue_ke bidang dari nilai pesan, nama bidang di sini adalahvalue_ts.CREATE TEMPORARY TABLE tempTable ( value_name STRING, value_ts STRING ) LIKE `kafkaJsonCatalog`.`kafka`.`testTopic`; -
-
Parameter tabel default
Parameter
Deskripsi
Keterangan
connector
Jenis konektor.
Nilainya adalah
kafkaatauupsert-kafka.topic
Nama topik yang sesuai.
Sama dengan nama tabel yang dideklarasikan.
properties.bootstrap.servers
Alamat broker Kafka.
Sama dengan pengaturan properties.bootstrap.servers katalog.
value.format
Format yang digunakan oleh konektor Kafka Flink untuk serialisasi atau deserialisasi nilai pesan Kafka.
Nilainya adalah
json.value.fields-prefix
Menentukan awalan kustom untuk semua bidang dari nilai pesan Kafka untuk menghindari konflik penamaan dengan bidang dari kunci pesan atau metadata.
Sama dengan pengaturan
value.fields-prefixkatalog.value.json.infer-schema.flatten-nested-columns.enable
Menentukan apakah akan memperluas kolom bersarang secara rekursif dalam JSON nilai pesan.
Sama dengan pengaturan
infer-schema.flatten-nested-columns.enablekatalog.value.json.infer-schema.primitive-as-string
Menentukan apakah akan melakukan inferensi semua tipe primitif dalam nilai pesan sebagai tipe STRING.
Sama dengan pengaturan
infer-schema.primitive-as-stringkatalog.value.fields-include
Menentukan cara nilai pesan menangani bidang dari kunci pesan.
Nilainya adalah
EXCEPT_KEY, yang berarti nilai pesan tidak mencakup bidang dari kunci pesan.Parameter ini dikonfigurasi hanya jika kunci pesan tidak kosong.
key.format
Format yang digunakan oleh konektor Kafka Flink untuk serialisasi atau deserialisasi kunci pesan Kafka.
Nilainya adalah
jsonatauraw.Parameter ini dikonfigurasi hanya jika kunci pesan tidak kosong.
Ketika kunci pesan tidak kosong tetapi penguraian gagal, nilainya adalah
raw. Jika penguraian berhasil, nilainya adalahjson.key.fields-prefix
Menentukan awalan kustom untuk semua bidang dari kunci pesan Kafka untuk menghindari konflik penamaan dengan bidang dari nilai pesan.
Sama dengan pengaturan
key.fields-prefixkatalog.Parameter ini dikonfigurasi hanya jika kunci pesan tidak kosong.
key.fields
Bidang yang menyimpan data yang diurai dari kunci pesan Kafka.
Daftar bidang kunci yang diurai diisi secara otomatis.
Parameter ini dikonfigurasi hanya jika kunci pesan tidak kosong dan tabel bukan tabel Upsert Kafka.