All Products
Search
Document Center

Realtime Compute for Apache Flink:Katalog JSON Kafka

Last Updated:Aug 20, 2026

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.

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.

    Catatan

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

    Catatan

    Dalam 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

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

    Penting

    Ganti 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'.

    Catatan

    Nilai 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'.

    Catatan

    Nilai 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'.

      Catatan

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

    Catatan

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

    Catatan

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

    Catatan

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

    Catatan

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

    Catatan

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

    Catatan

    Parameter ini hanya didukung di VVR 6.0.2 dan yang lebih baru.

  2. 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>'
    );
  3. Verifikasi bahwa katalog muncul di area Catalogs di sebelah kiri.

Katalog JSON Kafka

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

  2. 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), dan value_order (BIGINT), serta bidang virtual metadata partition (INT, primary key), offset (BIGINT, primary key), dan timestamp (TIMESTAMP_LTZ(3)). Untuk bidang virtual metadata ini, kolom extras menampilkan METADATA 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')*/;
    Catatan

    Untuk 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.示例

Catatan

Untuk contoh end-to-end penggunaan katalog JSON Kafka, lihat warehousing log waktu nyata.

Hapus katalog JSON Kafka

Peringatan

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.

  1. Masukkan perintah berikut di editor SQL pada halaman Scripts.

    DROP CATALOG ${catalog_name};

    Ganti ${catalog_name} dengan nama katalog JSON Kafka Anda.

  2. Pilih perintah DROP CATALOG, klik kanan, lalu pilih Run.

  3. 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.Schema合并

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

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

    Catatan

    Jika 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 ts dengan format '2023-01-01 12:00:01', Katalog JSON Kafka secara otomatis melakukan inferensi bidang ts sebagai tipe TIMESTAMP. Untuk menggunakan bidang ts sebagai tipe STRING, deklarasikan tabel sebagai berikut. Karena konfigurasi default menambahkan awalan value_ ke bidang dari nilai pesan, nama bidang di sini adalah value_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 kafka atau upsert-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-prefix katalog.

    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.enable katalog.

    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-string katalog.

    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 json atau raw.

    Parameter ini dikonfigurasi hanya jika kunci pesan tidak kosong.

    Ketika kunci pesan tidak kosong tetapi penguraian gagal, nilainya adalah raw. Jika penguraian berhasil, nilainya adalah json.

    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-prefix katalog.

    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.