All Products
Search
Document Center

Realtime Compute for Apache Flink:Konektor Kafka YAML

Last Updated:Sep 16, 2026

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.

    Penting

    Batasan 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.idempotence secara 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 menambahkan properties.enable.idempotence=false ke 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-prefix yang 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

    Catatan

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

    Catatan

    Anda 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 dengan user_event_

    • prod\.logs\..*: mencocokkan topik yang namanya memiliki awalan prod.logs. (. harus di-escape)

    Catatan

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

    Catatan

    Parameter 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:300

    scan.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.id diduplikasi.

    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

    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.

    Catatan

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

    Catatan

    Nilai 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 (,).

    Catatan

    Kolom 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, gunakan CREATE 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.

    Catatan

    Parameter 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

    Catatan

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

    Catatan

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

    Catatan

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

      Catatan

      Parameter ini hanya didukung pada VVR 8.0.11 dan yang lebih baru.

      Penting

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

      Catatan

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

      Catatan

      Parameter ini hanya didukung pada VVR 8.0.11 dan yang lebih baru.

      Penting

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

      Catatan

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

      Catatan

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

      Catatan

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

      Catatan

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

    Catatan

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

    Catatan

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

      Catatan

      Parameter 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

      db

      Schema

      schema

      Tabel

      table

      Pemetaan saat parameter ini dinonaktifkan:

      Bagian ID Tabel CDC

      Kunci Debezium JSON

      Namespace

      None

      Schema

      db

      Tabel

      table

      Catatan

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

Catatan

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:

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

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

Catatan

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.

Penting

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.

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

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

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

      image

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

  1. Implementasikan antarmuka KafkaPayloadConverter dan paketkan. Repositori demo open source sudah menyediakan beberapa contoh implementasi.

  2. Tambahkan file yang telah dipaketkan ke dependensi tambahan pekerjaan ingesti data.

  3. 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
  4. 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