全部產品
Search
文件中心

Realtime Compute for Apache Flink:Kafka資料攝入連接器

更新時間:Sep 16, 2026

Kafka連接器可以用於Flink CDC資料攝入作業開發,作為源端讀取或目標端寫入。本文介紹Kafka 資料攝入連接器的文法結構、參數配置和使用樣本。

前提條件

請根據需求選擇以下任意一種方式串連叢集:

  • 串連阿里雲雲訊息佇列Kafka版叢集

    • Kafka叢集版本在0.11及以上。

    • ApsaraMQ for Kafka叢集已建立。詳情請參見建立資源。

    • Flink工作空間與Kafka叢集處於同一VPC內,且ApsaraMQ for Kafka已對Flink開放白名單,具體操作請參見配置白名單。

    重要

    寫入阿里雲Kafka的限制:

    • 阿里雲Kafka不支援zstd壓縮格式寫入。

    • 阿里雲Kafka不支援等冪和事務寫入,無法使用Kafka結果表提供的精確一次語義exactly-once semantic功能。自Realtime Compute引擎VVR 8.0.0版本起,Kafka Connector 使用的開源Kafka Client版本升級到3.x版本,該Client的properties.enable.idempotence屬性預設值從false變為true表示顯式開啟等冪寫入, 因此在使用Realtime Compute引擎VVR 8.0.0及以上版本寫入阿里雲Kafka時,需要在結果表中添加顯式配置properties.enable.idempotence=false以關閉等冪寫入功能,避免無法寫入阿里雲Kafka的問題。阿里雲Kafka的儲存引擎對比與功能限制參見儲存引擎對比。

  • 串連自建Apache Kafka叢集

    • 自建Apache Kafka叢集版本在0.11及以上。

    • Flink與自建Apache Kafka叢集之間的網路已打通。如何通過公網串連自建叢集,詳情請參見網路連接選型。

    • 僅支援Apache Kafka 2.8版本的用戶端配置項,詳情請參見Apache Kafka消費者和生產者配置項文檔。

使用限制

  • 建議在Realtime Compute引擎VVR 11.1及以上版本使用Kafka作為Flink CDC資料攝入的同步資料來源。

  • 僅支援JSON、Debezium JSON和Canal JSON格式,其他資料格式暫不支援。

  • 對於資料來源,僅Realtime Compute引擎VVR 8.0.11及以上版本支援同一張表的資料分布在多個分區。

注意事項

目前不推薦使用事務寫入,這是 Flink 社區和 Kafka 社區的設計缺陷所致。當設定sink.delivery-guarantee = exactly-once,Kafka Connector 會啟用事務寫入,存在三個已知問題:

  • 每個 Checkpoint 會產生一個 Transaction ID。如果 Checkpoint 間隔太短,Transaction ID會過多。Kafka 叢集的 Coordinator 可能因此記憶體不足,從而破壞 Kafka 叢集的穩定性。

  • 每個事務會建立一個 Producer 執行個體。如果同時提交的事務太多,TaskManager 的記憶體可能耗盡,從而破壞 Flink 作業的穩定性。

  • 多個 Flink 作業若使用相同的sink.transactional-id-prefix,它們產生的事務 ID 可能衝突。一個作業寫入失敗時,會阻塞 Kafka 分區的 LSO(Log Start Offset)前進,這會影響所有消費者讀取該分區的資料。

如果你需要 Exactly-Once 語義,改用 Upsert Kafka 寫入主鍵表,並用主鍵保證等冪性。如果需要使用事務寫入,請參見Kafka SQL連接器。

文法結構

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

配置項

  • 通用

    參數

    說明

    是否必填

    資料類型

    預設值

    備忘

    type

    源端或目標端類型。

    是

    String

    無

    固定值為kafka

    name

    源端或目標端名稱。

    否

    String

    無

    無

    properties.bootstrap.servers

    Kafka broker地址。

    是

    String

    無

    格式為host:port,host:port,host:port,以英文逗號(,)分割。

    properties.*

    對Kafka用戶端的直接配置。

    否

    String

    無

    尾碼名必須是Kafka官方文檔中定義的生產者和消費者配置。

    Flink會將properties.首碼移除,並將剩餘的配置傳遞給Kafka用戶端。例如可以通過'properties.allow.auto.create.topics' = 'false' 來禁用自動建立topic。

    key.format

    讀取或寫入Kafka訊息key部分使用的格式。

    否

    String

    無

    • 對Source,僅支援json。

    • 對Sink,取值如下:

      • csv

      • json

    說明

    僅Realtime Compute引擎VVR 11.0.0及以上版本支援該參數。

    value.format

    讀取或寫入Kafka訊息value部分時使用的格式。

    否

    String

    debezium-json

    • 對Source,取值如下:

      • debezium-json 

      • canal-json

      • json

    • 對Sink,取值如下:

      • debezium-json 

      • canal-json

      • canal-protobuf

    說明
    • 僅Realtime Compute引擎VVR 8.0.10及以上版本支援debezium-json和canal-json格式。

    • 僅Realtime Compute引擎VVR 11.0.0及以上版本支援json格式。

  • 源表

    參數

    說明

    是否必填

    資料類型

    預設值

    備忘

    topic

    讀取的topic名稱。

    否

    String

    無

    以英文分號 (;) 分隔多個topic名稱,例如topic-1和topic-2

    說明

    topic和topic-pattern兩個選項只能指定其中一個。

    topic-pattern

    匹配讀取topic名稱的Regex。所有匹配該Regex的topic在作業運行時均會被讀取。

    否

    String

    無

    樣本:

    • user_event_.*:匹配所有以 user_event_ 開頭的 topic

    • prod\.logs\..*:匹配 prod.logs. 首碼的 topic(. 需轉義)

    說明

    topic和topic-pattern兩個選項只能指定其中一個。

    properties.group.id

    消費組ID。

    否

    String

    無

    如果指定的group id為首次使用,則必須將properties.auto.offset.reset設定為earliest或latest以指定初次開機位點。

    scan.startup.mode

    Kafka讀取資料的啟動位點。

    否

    String

    group-offsets

    取值如下:

    • earliest-offset:從Kafka最早分區開始讀取。

    • latest-offset:從Kafka最新位點開始讀取。

    • group-offsets(預設值):從指定的properties.group.id已提交的位點開始讀取。

    • timestamp:從scan.startup.timestamp-millis指定的時間戳記開始讀取。

    • specific-offsets:從scan.startup.specific-offsets指定的位移量開始讀取。

    說明

    該參數在作業無狀態啟動時生效。作業在從checkpoint重啟或狀態恢複時,會優先使用狀態中儲存的進度恢複讀取。

    scan.startup.specific-offsets

    specific-offsets啟動模式下,指定每個分區的啟動位移量。

    否

    String

    無

    例如partition:0,offset:42;partition:1,offset:300

    scan.startup.timestamp-millis

    timestamp啟動模式下,指定啟動位點時間戳記。

    否

    Long

    無

    單位為毫秒

    scan.topic-partition-discovery.interval

    動態檢測Kafka topic和partition的時間間隔。

    否

    Duration

    5分鐘

    分區檢查間隔預設為5分鐘。需要顯式地設定分區檢查間隔為非正數才能關閉此功能。開啟動態分區發現後,Kafka Source 可以自動地發現新增的分區並自動讀取對應分區上的資料。在topic-pattern模式下,不僅讀取已有topic的新增分區資料,也會讀取符合正則匹配的新增topic的所有分區資料。

    scan.check.duplicated.group.id

    是否檢查通過properties.group.id指定的消費者組有重複。

    否

    Boolean

    false

    參數取值如下:

    • true:在啟動作業前檢查消費者組是否有重複,如有重複作業將會報錯,避免與現有的消費者組產生衝突。

    • false:直接啟動作業,不檢查消費者組衝突。

    schema.inference.strategy

    Schema解析策略。

    否

    String

    continuous

    取值如下:

    • continuous:對每條資料均進行 Schema 解析。在前後 Schema 不相容時,解析出更寬的 Schema 併產生 Schema 變更事件。

    • static:僅在作業啟動時進行一次Schema解析,後續根據初始Schema解析資料,不會產生Schema變更事件。

    說明
    • Schema解析詳情可見Kafka SQL連接器。

    • 僅VVR 8.0.11及以上版本支援該配置項。

    scan.max.pre.fetch.records

    Schema初始解析時,對每個分區最多嘗試消費解析的訊息數量

    否

    Int

    50

    在作業實際讀取並處理資料前,對每個分區嘗試提前消費指定數量的最新訊息,用於初始化Schema資訊。

    key.fields-prefix

    自訂添加到訊息鍵(Key)解析出欄位名稱的首碼,以避免Kafka訊息鍵解析後的命名衝突問題。

    否

    String

    無

    假設該配置項設為key_,當key中包含欄位名a時,解析key後該欄位名稱為key_a。

    說明

    key.fields-prefix的配置值不可以是value.fields-prefix的首碼。

    value.fields-prefix

    自訂添加到訊息體(Value)解析出欄位名稱的首碼,以避免Kafka訊息體解析後的命名衝突問題。

    否

    String

    無

    假設該配置項設為value_,當value中包含欄位名b時,解析value後該欄位名稱為value_b。

    說明

    value.fields-prefix的配置值不可以是key.fields-prefix的首碼。

    metadata.list

    需要傳遞給下遊的中繼資料列

    否

    String

    無

    可用的中繼資料列包括topic、partition、offset、timestamp、timestamp-type、headers、leader-epoch、__raw_key__和__raw_value__。使用英文逗號分隔。

    說明

    __raw_key__和 __raw_value__中繼資料列需要VVR 11.6及更高版本可用。

    scan.value.initial-schemas.ddls

    通過DDL方式指定某些表的初始Schema。

    否

    String

    無

    多個DDL使用英文;串連。例如,使用CREATE TABLE db1.t1 (id BIGINT, name VARCHAR(10)); CREATE TABLE db1.t2 (id BIGINT);為表db1.t1和表db1.t2分別指定初始Schema。

    這裡的DDL表結構應當和寫入目標表保持一致,並且滿足Flink SQL的文法規則。

    說明

    VVR 11.5及以上版本支援該配置。

    ingestion.ignore-errors

    是否忽略資料解析過程中的報錯。

    否

    Boolean

    false

    說明

    VVR 11.5及以上版本支援該配置。

    ingestion.error-tolerance.max-count

    忽略資料解析過程報錯時,累計多少次報錯後作業失敗。

    否

    Integer

    -1

    僅當ingestion.ignore-errors開啟時生效,預設值-1表示解析異常不觸發作業失敗。

    說明

    VVR 11.5及以上版本支援該配置。

    scan.duplicate-field.strategy

    在 Key 和 Value 部分解析出重複欄位名時,應該如何進行處理。

    否

    String

    EXCEPTION

    參數取值如下:

    • EXCEPTION:在 key 和 value 存在重複欄位時,拋出異常。這是 VVR 11.6 及更早版本的預設行為。

    • PREFER_KEY:在欄位重複時,優先選擇 key 欄位的值。

    • PREFER_VALUE:在欄位重複時,優先選擇 value 欄位的值。

    說明

    VVR 11.7 及以上版本支援該配置。

    • 源表 Debezium JSON 格式

      參數

      是否必填

      資料類型

      預設值

      描述

      debezium-json.distributed-tables

      否

      Boolean

      false

      如果Debezium JSON內單張表資料會出現在多個分區,則需要開啟此選項。

      說明

      僅VVR 8.0.11及以上版本支援該配置項。

      重要

      修改該配置項後,需要無狀態啟動作業。

      debezium-json.schema-include

      否

      Boolean

      false

      設定Debezium Kafka Connect時,可以啟用Kafka配置value.converter.schemas.enable,以在訊息中包含schema。此選項表明Debezium JSON訊息是否包含schema。

      參數取值如下:

      • true:Debezium JSON訊息包含schema。

      • false:Debezium JSON訊息不包含schema。

      debezium-json.ignore-parse-errors

      否

      Boolean

      false

      參數取值如下:

      • true:當解析異常時,跳過當前行。

      • false(預設值):報出錯誤,作業啟動失敗。

      debezium-json.infer-schema.primitive-as-string

      否

      Boolean

      false

      解析表結構時,是否解析所有類型為String類型。

      參數取值如下:

      • true:解析所有基本類型為String。

      • false(預設值):按照基本規則進行解析。

      debezium-json.infer-schema.string-type-inference

      否

      Boolean

      true

      是否嘗試將字串欄位推斷為 TIME、DATE 或 TIMESTAMP類型。設定為 false 時跳過該推斷,欄位保持 STRING。

      說明

      VVR 11.8及以上版本支援該配置。

    • 源表 Canal JSON 格式

      參數

      是否必填

      資料類型

      預設值

      描述

      canal-json.distributed-tables

      否

      Boolean

      false

      如果Canal JSON內單張表資料會出現在多個分區,則需要開啟此選項。

      說明

      僅VVR 8.0.11及以上版本支援該配置項。

      重要

      修改該配置項後,需要無狀態啟動作業。

      canal-json.database.include

      否

      String

      無

      一個可選的Regex,通過正則匹配Canal記錄中的database元欄位,僅讀取指定資料庫的changelog記錄。正則字串與Java的Pattern相容。

      canal-json.table.include

      否

      String

      無

      一個可選的Regex,通過正則匹配Canal記錄中的table元欄位,僅讀取指定表的changelog記錄。正則字串與Java的Pattern相容。

      canal-json.ignore-parse-errors

      否

      Boolean

      false

      參數取值如下:

      • true:當解析異常時,跳過當前行。

      • false(預設值):報出錯誤,作業啟動失敗。

      canal-json.infer-schema.primitive-as-string

      否

      Boolean

      false

      解析表結構時,是否所有類型解析為String類型。

      參數取值如下:

      • true:解析所有基本類型為String。

      • false(預設值):按照基本規則進行解析。

      canal-json.infer-schema.strategy

      否

      String

      AUTO

      解析表結構時的解析策略。

      參數取值如下:

      • AUTO(預設值):通過解析JSON資料自動解析。如果資料中不包含sqlType欄位,建議使用AUTO以避免解析失敗。

      • SQL_TYPE:通過canal json資料中的sqlType數組解析。如果資料包含sqlType欄位時,建議將canal-json.infer-schema.strategy設定為SQL_TYPE以獲得更精確的類型。

      • MYSQL_TYPE:通過canal json資料中的mysqlType數組解析。

      當Kafka中的Canal JSON資料包含sqlType欄位且需要更精確的類型映射時,建議將canal-json.infer-schema.strategy設定為SQL_TYPE。

      sqlType類型映射規則請見Kafka SQL連接器。

      說明
      • VVR 11.1及以上版本支援該配置。

      • MYSQL_TYPE在VVR 11.3及以上版本支援使用。

      canal-json.mysql.treat-mysql-timestamp-as-datetime-enabled

      否

      Boolean

      true

      是否將mysql timestamp類型映射到cdc timestamp類型:

      • true(預設):mysql timestamp類型映射到cdc timestamp類型。

      • false:mysql timestamp類型映射到cdc timestamp_ltz類型。

      canal-json.mysql.treat-tinyint1-as-boolean.enabled

      否

      Boolean

      true

      使用MYSQL_TYPE解析時,是否將mysql tinyint(1)類型映射到cdc boolean類型:

      • true(預設):mysql tinyint(1)類型映射到cdc boolean類型。

      • false:mysql tinyint(1)類型映射到cdc tinyint(1)類型。

      該配置僅當canal-json.infer-schema.strategy配置為MYSQL_TYPE時生效。

      canal-json.infer-schema.string-type-inference

      否

      Boolean

      true

      是否嘗試將字串欄位推斷為 TIME、DATE 或 TIMESTAMP類型。設定為 false 時跳過該推斷,欄位保持 STRING。

      說明

      VVR 11.8及以上版本支援該配置。

    • 源表 JSON 格式

      參數

      是否必填

      資料類型

      預設值

      描述

      json.timestamp-format.standard

      否

      String

      SQL

      指定輸入和輸出時間戳記格式。參數取值如下:

      • SQL:解析yyyy-MM-dd HH:mm:ss.s{precision}格式的輸入時間戳記,例如2020-12-30 12:13:14.123。

      • ISO-8601:解析yyyy-MM-ddTHH:mm:ss.s{precision}格式的輸入時間戳記,例如2020-12-30T12:13:14.123。

      json.ignore-parse-errors

      否

      Boolean

      false

      參數取值如下:

      • true:當解析異常時,跳過當前行。

      • false(預設值):報出錯誤,作業啟動失敗。

      json.infer-schema.primitive-as-string

      否

      Boolean

      false

      解析表結構時,是否解析所有類型為String類型。

      參數取值如下:

      • true:解析所有基本類型為String。

      • false(預設值):按照基本規則進行解析。

      json.infer-schema.flatten-nested-columns.enable

      否

      Boolean

      false

      解析JSON格式資料時,是否遞迴式地展開JSON中的嵌套列。參數取值如下:

      • true:遞迴式展開。

      • false(預設值):將嵌套列當作String處理。

      json.decode.parser-table-id.fields

      否

      String

      無

      解析JSON格式資料時,是否使用部分JSON欄位值產生tableId,多個欄位使用英文,串連。例如:JSON資料為{"col0":"a", "col1","b", "col2","c"},產生結果如下:

      配置

      tableId

      col0

      a

      col0,col1

      a.b

      col0,col1,col2

      a.b.c

      json.infer-schema.fixed-types

      否

      String

      無

      解析JSON格式資料時,指定某些欄位的具體類型,多個欄位使用英文,串連。例如:id BIGINT, name VARCHAR(10)可以指定JSON資料中的id欄位類型為BIGINT,name欄位類型為VARCHAR(10)。

      說明
      • VVR 11.5及以上版本支援該配置。

      • VVR 11.5版本使用該配置時,需要額外添加配置scan.max.pre.fetch.records: 0。

      json.decode.converter-class

      否

      String

      無

      在解析JSON資料前,通過converter可以對JSON byte[]提前進行修改,填寫具體實作類別的全限定名。

      說明

      VVR 11.6及以上版本支援該配置。

      json.decode.empty-value-as-delete.enabled

      否

      Boolean

      false

      是否將Kafka compacted topic中的墓碑訊息(value 為空白)解析為 DELETE 事件,用於 compacted topic 鏡像、CDC 刪除訊號等以空 value 表示刪除語義的情境。

      說明

      VVR 11.7 及以上版本支援該配置。

      json.infer-schema.string-type-inference

      否

      Boolean

      true

      是否嘗試將字串欄位推斷為 TIME、DATE 或 TIMESTAMP類型。設定為 false 時跳過該推斷,欄位保持 STRING。

      說明

      VVR 11.8及以上版本支援該配置。

  • 結果表

    參數

    說明

    是否必填

    資料類型

    預設值

    備忘

    type

    目標端類型。

    是

    String

    無

    固定值為kafka

    name

    目標端名稱。

    否

    String

    無

    無

    topic

    Kafka Topic名稱。

    否

    String

    無

    開啟時,所有的資料都會寫入這個Topic。

    說明

    如果沒有開啟,每條資料會寫入到其TableID對應字串(通過.拼接產生)的Topic,例如databaseName.tableName。

    partition.strategy

    資料寫入Kafka分區的策略。

    否

    String

    all-to-zero

    取值如下:

    • all-to-zero(預設值):將所有資料寫入 0 號分區。

    • hash-by-key:根據主鍵的雜湊值將資料寫到多個分區。保證同一個主鍵的資料在同一個分區並且有序。

    sink.tableId-to-topic.mapping

    上遊表名到下遊 Kafka Topic名的映射關係。 

    否

    String

    無

    每個映射關係由;分割,上遊表的表名和下遊 Kafka 的Topic名由:分割,表名可以使用Regex,映射到同一個topic的多張表可以使用,拼接。例如mydb.mytable1:topic1;mydb.mytable2:topic2。

    說明

    配置這個參數能夠在保留原始表名資訊的同時修改映射的Topic。

    • 結果表Canal JSON格式

      參數

      是否必填

      資料類型

      預設值

      描述

      canal-json.serialize.update.keep-changed-fields-only

      否

      Boolean

      false

      寫入 canal-json 格式的 UPDATE 訊息時,old 部分是否僅包含發生變化欄位的變更前的值。

      說明

      VVR 11.8及以上版本支援該配置。

    • 結果表Debezium JSON格式

      參數

      是否必填

      資料類型

      預設值

      描述

      debezium-json.include-schema.enabled

      否

      Boolean

      false

      Debezium JSON資料中是否包含Schema資訊。

      debezium-json.emit.full-table-id.enabled

      否

      Boolean

      false

      是否將完整的三段式Table ID寫入Debezium JSON中繼資料欄位中。

      啟用此參數時的映射關係為:

      CDC Table ID部分

      Debezium JSON鍵

      Namespace

      db

      Schema

      schema

      Table

      table

      禁用此參數時的映射關係為:

      CDC Table ID部分

      Debezium JSON鍵

      Namespace

      無

      Schema

      db

      Table

      table

      說明

      VVR 11.6及以上版本支援此配置。

複用已有 Catalog

自VVR 11.5版本起,您可以在Flink CDC資料攝入作業中直接引用“資料管理”頁面中建立的內建Kafka Catalog,減少手寫串連屬性工作量。

source:
  type: kafka
  using.built-in-catalog: kafka_catalog

目前,資料攝入作業支援自動複用以下 Kafka Catalog 參數:

  • 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

如果希望覆蓋以上自動複用的參數,可顯式寫出相應的 YAML 參數,其具備更高的優先順序。

配置樣本

  • 使用 Kafka 作為資料攝入源端:

    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
  • 使用 Kafka 作為資料攝入目標端:

    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}

    其中,使用route模組以設定源表寫入Kafka的Topic名稱。

說明

阿里雲Kafka預設不開啟自動建立Topic功能,參見自動化建立Topic相關問題,寫入到阿里雲Kafka時,需要預先建立對應的Topic,詳情請參見步驟三:建立資源。

使用樣本

下面展示幾個典型使用情境的配置樣本。

讀取單獨topic

讀取topic customers,寫入阿里雲資料湖的配置樣本:

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

其中,Json格式預設產生的表名和topic名稱相同。

讀取多個topic

讀取符合Regex的多個topic,寫入 StarRocks SQL連接器 的配置樣本:

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,
  # 否則舊主鍵對應的行會殘留在下遊

其中,Json格式預設產生的表名和topic名稱相同。

讀取key並防止欄位衝突

為了防止 key 和 value 的欄位衝突導致報錯,有以下解決方案:

  1. 為欄位添加首碼,避免欄位衝突:

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. 配置衝突解決方案策略,詳情請見Kafka SQL連接器。如下配置會優先使用key中欄位,出現重複時忽略value中的同名欄位。

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

增加中繼資料列

讀取topic customers,並在欄位增加中繼資料topic和partition的配置樣本:

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

解析報錯處理

資料解析報錯會導致作業失敗,支援配置容忍解析報錯,常配合髒資料收集使用。

完全忽略解析報錯的配置樣本:

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

累計解析失敗30次後作業報錯的配置樣本:

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

JSON資料讀取

接下來介紹 Json 格式讀取常見的一些使用方式。

指定 table id 解析

預設 json 資料的 table id 是 topic 名稱,可以指定資料中的欄位值作為table id。如下樣本將欄位 db 和 tbl 指定為 table id。

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

指定欄位類型

欄位類型解析按照欄位值推導擷取,有些欄位類型的推斷並不是使用者期望的類型。支援指定某些欄位的固定類型,並跳過後續對這個欄位類型的推導和演變。

如下配置固定了4個欄位的類型:

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

Canal Json資料讀取

接下來介紹 Canal Json 格式讀取常見的一些使用方式。

類型推導策略

預設情況下,讀取 Canal json 資料會使用資料值進行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

此外支援使用 canal json 資料中記錄的 Schema 資訊(sql type 或 mysql type)進行推導,如下例子使用 mysql type 推導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

Debezium Json資料讀取

讀取topic customers,寫入阿里雲資料湖的配置樣本:

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

MySQL 原始 binlog 同步 Kafka

Flink CDC資料攝入支援同步MySQL 原始binlog內容到 canal json,如下作業可以同步多張表的binlog到 topic 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

進階功能

表結構解析和變更同步策略

Kafka連接器會維護當前已知的所有表的Schema。

表結構資訊初始化

表結構資訊包含欄位和資料類型資訊、庫表資訊和主鍵資訊,這三種資訊的初始化方式如下:

  • 欄位和資料類型資訊

資料攝入作業能夠根據資料自動推匯出表的欄位和資料類型資訊,然而在某些情境下,您可能希望指定某些表的欄位和類型,根據使用者指定欄位類型的粒度,表結構資訊的初始化有以下三種配置策略:

  1. 完全由程式推導表結構

在讀取Kafka資料前,Kafka連接器會預先在每個分區中嘗試消費最多scan.max.pre.fetch.records條訊息,解析每條資料的Schema,再將這些Schema合并,用於初始化表結構資訊。後續在實際消費資料前會根據初始化的Schema產生對應的建表事件。

說明

對於Debezium JSON和Canal JSON格式,表資訊在具體訊息中,提前消費的scan.max.pre.fetch.records條訊息中可能包含了若干個表的資料,因此對每張表而言,提前消費的資料條數無法確定。預消費和初始化表結構資訊只會在實際消費和處理每個分區的訊息前進行一次,若後續有新表資料,該表的第一條資料解析出的表結構會作為初始表結構,不會重新預消費和初始化對應的表結構。

重要

僅VVR 8.0.11及以上版本支援單表資料分布在多個分區,對於該情境需要將配置項debezium-json.distributed-tables或canal-json.distributed-tables設為true。

  1. 指定初始表結構

在某些情境下,您希望自行指定初始的表結構,例如將 Kafka 的資料寫入到一張預先建立的下遊表中。這時您可以通過添加scan.value.initial-schemas.ddls參數指定初始的表結構。配置樣本如下:

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);

這裡的建表語句需要和目標表的表結構保持一致。這樣,我們就為 db1.t1 這張表指定了 id 這個欄位的初始類型為 BIGINT,name 這個欄位的初始類型為 VARCHAR(10) ,為 db1.t2 這張表指定了 id 這個欄位的初始類型為 BIGINT。

這裡的建表語句使用的是 FlinkSQL 的文法。

  1. 指定欄位設定為固定類型

在某些情境下,您希望固定特定欄位的資料類型,例如對於某些可能被推導為 TIMESTAMP 類型的欄位,您希望以字串的格式下發,這時您可以通過添加json.infer-schema.fixed-types參數指定初始的表結構(僅在訊息格式為 json 時有效)。配置樣本如下:

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

這樣,我們就為所有的 id 欄位指定了類型固定為 BIGINT,所有的 name 欄位指定了類型固定為 VARCHAR(10)。

這裡的類型與 FlinkSQL 的資料類型一致。

  • 庫表資訊

    • 對於Canal JSON和Debezium JSON格式,表資訊由具體訊息解析得到,包括資料庫和表名。

    • 對於JSON格式,預設配置下,表資訊僅包含表名,即資料所在的topic名稱。如果您的資料中包含庫表資訊,可以通過json.infer-schema.fixed-types指定庫表資訊所在的欄位,我們會將這些欄位對應到庫名/表名中。配置樣本如下:

      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

      這樣,我們會將每條資料下發到庫名為col1欄位的值,表名為col2欄位的值的表中。

  • 主鍵資訊

    • 對於Canal JSON格式,會根據JSON中的pkNames欄位定義表的主鍵。

    • 對於Debezium JSON和JSON格式,JSON中不包含主鍵資訊,可以通過transform規則手動為表添加主鍵:

      transform:
        - source-table: \.*.\.*
          projection: \*
          primary-keys: key1, key2

Schema解析和Schema變更

在表結構初始化完成後,若schema.inference.strategy配置為static,Kafka連接器會根據初始的表結構解析每個訊息的訊息體(Value),不會產生Schema變更事件。若schema.inference.strategy配置為continuous,Kafka連接器會解析每個Kafka訊息的訊息體,解析出訊息的物理列,並與當前維護的Schema比對,若解析出的Schema與當前Schema不一致時,會嘗試將Schema合并,同時產生對應的表結構變更事件,合并規則如下:

  • 如果解析出的物理列中包含當前Schema中沒有的欄位,則會將這些欄位加入到Schema中,同時產生新增可空列事件。

  • 如果解析出的物理列中不包含當前Schema中已有的欄位,該欄位仍會保留,該列的資料會填充為NULL,不產生刪除列事件。

  • 如果兩者出現了同名列,則按照以下情境進行處理:

    • 當類型相同且精度不同時,會取兩者中較大精度的類型,同時產生列類型變更事件。

    • 當類型不同時,會按照如下圖的樹形結構找到最小父節點,作為該同名列的類型,同時產生列類型變更事件。

      image

  • 當前支援的Schema變更策略如下:

    • 添加列:會在當前Schema末尾添加對應的列,並同步新增列的資料,新增的列會設定為可空列。

    • 刪除列:不會產生刪除列事件,而是後續將該列的資料自動填滿為NULL值。

    • 重新命名列:被看作為添加列和刪除列,在當前Schema末尾添加重新命名後的列,並將重新命名前的列資料填充為NULL值。

    • 列類型變更:

      • 對於支援列類型變更的下遊系統,在下遊Sink支援處理列類型變更後,資料攝入作業支援普通列的類型變更,例如,從INT類型變更到BIGINT類型。此類變更依賴於下遊Sink支援的列類型變更規則,不同的結果表支援的列類型變更規則也不相同,請參考結果表文檔擷取其支援的列類型變更規則。

      • 對於不支援列類型變更的下遊系統,比如Hologres,此類情境可以使用寬類型映射,即作業啟動時在下遊系統建立類型更加寬泛的表,在列類型變更發生時判斷該類型變更下遊Sink是否可以接受從而實現寬容的列類型變更支援。

  • 當前暫不支援的Schema變更:

    • 主鍵或索引等約束的變更。

    • 從NOT NULL轉為NULLABLE變更。

  • Canal JSON的Schema解析

    Canal JSON資料中可能包含可選的sqlType欄位,其中記錄了資料列的精確類型資訊。為了擷取更準確的Schema,可以通過將canal-json.infer-schema.strategy配置為SQL_TYPE使用sqlType中的類型。類型映射關係如下:

    JDBC類型

    Type Code

    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

    其他類型

髒資料容忍與收集

在某些情況下,您的 Kafka 資料來源中可能包含一些格式不正確的資料(髒資料),為了避免同步作業因為這些髒資料導致頻繁失敗重啟,您可以配置忽略這部分異常資料。配置樣本如下:

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

這裡配置了忽略 1000 條髒資料的策略,讓您的作業能夠在存在小批量髒資料時正常運行,而在髒資料條數超過這個閾值時讓作業進入失敗狀態,提醒您對資料進行校正。

如果您希望作業永遠不因為髒資料的存在導致失敗,可以採用如下的配置:

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

髒資料容忍策略保障了作業不會因為異常資料頻繁失敗,您可能還希望進一步的瞭解這些髒資料的資訊以調整 Kafka 資料生產者的行為,參考髒資料收集中介紹的流程,您可以在 TaskManager 日誌中查看到作業的髒資料,配置樣本如下:

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

表名與Topic的映射策略

在使用kafka作為資料攝入作業的目標端時,由於寫入到Kafka訊息格式(debezium-json或者canal-json)中還包含表名資訊,後續消費Kafka訊息時往往以資料中的表名資訊作為實際表名(而非topic名稱),因此需要謹慎配置表名與Topic的映射策略。

假設在MySQL有mydb.mytable1,mydb.mytable2兩張表需要同步,可能的配置策略有以下幾種:

1. 不配置任何映射策略

在沒有任何映射策略的情況下,每張表會寫入到對應的由“庫名.表名”組成的topic中。因此mydb.mytable1的資料會寫入到名為mydb.mytable1的topic中,mydb.mytable2的資料會寫入到名為mydb.mytable2的topic中。配置樣本如下:

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. 配置route規則進行映射(不推薦)

在很多情境下,使用者不希望寫入的topic直接為“庫名.表名”的格式,希望將資料寫入到指定的topic中,因此會配置route規則進行映射。配置樣本如下:

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

此時所有來自mydb.mytable1,mydb.mytable2的資料都會寫入到mytable1這一個topic中。

然而,通過route規則修改寫入的topic名稱時,也會修改Kafka訊息(debezium-json或者canal-json格式)中的表名資訊,此時Kafka訊息中所有的表名都為mytable1,在其他系統消費這個topic的Kafka訊息時,可能出現不符合預期的情況。

3. 配置sink.tableId-to-topic.mapping參數進行映射(推薦)

為了在配置表名與Topic的映射規則的同時保留源表表名資訊,可以使用sink.tableId-to-topic.mapping參數完成該需求。配置樣本如下:

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

或者

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

此時所有來自mydb.mytable1,mydb.mytable2的資料都會寫入到mytable1這一個topic中,並且Kafka訊息(debezium-json或者canal-json格式)中的表名資訊仍然為mydb.mytable1或者mydb.mytable2,在其他系統消費這個topic的Kafka訊息時,能正確擷取源表表名資訊。

JSON Converter實現

JSON格式的資料有時無法保證格式完全一致,或者有處理需求無法通過Transform配置模組完成,例如合并兩個列為新的列並刪除原有的兩列,並保證能進行 Schema 變更同步。為了更靈活處理JSON資料,可以通過實現KafkaPayloadConverter介面在架構處理前,提前對JSON資料進行修改,使用時需要添加配置json.decode.converter-class填寫具體實作類別的全限定名。

實現和使用JSON Converter的流程如下:

  1. 實現KafkaPayloadConverter介面並打包,在開源Demo倉庫已經提供了一些實現樣本。

  2. 將打好的包添加到資料攝入作業的附加依賴。

  3. 在Source模組填寫json.decode.converter-class配置,以開源專案中的ArrayElementExtractorConverter類為例:

    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. 部署作業並運行。

key format和value format

在VVR 11.5及以前的版本中,key format 和 value format的配置無法進行區分,當使用相同的format時,format配置會同時應用在key format和value format。

VVR 11.6及以後的版本對這部分進行了最佳化,以json format 為例,format 的配置項的路由規則如下:

  • format首碼(如json.infer-schema.primitive-as-string): 為了保持版本相容,format首碼的配置會預設同時應用到 key format和 value format。

  • key首碼加format首碼(如key.json.infer-schema.primitive-as-string): key首碼加format首碼的配置會只應用在key format,優先順序比format首碼的配置更高。

  • value首碼加format首碼(如value.json.infer-schema.primitive-as-string): value首碼加format首碼的配置會只應用在value format,優先順序比format首碼的配置更高。

如下的Kafka source配置,為key format和value format配置了不同的JSON Converter。

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