全部產品
Search
文件中心

Data Lake Formation:Flink CDC訪問DLF

更新時間:Aug 25, 2026

如何在阿里雲Realtime ComputeFlink版上通過Flink CDC以Paimon REST訪問DLF Catalog。

前提條件

使用限制

僅Realtime Compute引擎VVR 11.1.0及以上版本支援對接DLF Catalog。

建立DLF Catalog

詳情請參見 DLF 快速入門

Flink CDC對接Catalog配置參數

建立資料攝入作業的操作流程,請參見Flink CDC資料攝入作業開發

Flink中資料攝入作業的Sink使用以下配置:

sink:
  type: paimon
  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

配置項說明如下:

配置項

描述

是否必填

樣本

catalog.properties.metastore

Metastore類型,固定為rest。

rest

catalog.properties.token.provider

Token提供方,固定為dlf。

dlf

catalog.properties.uri

訪問DLF Rest Catalog Server的URI,格式為http://[region-id]-vpc.dlf.aliyuncs.com。詳見地區與服務存取點中的Region ID。

http://cn-hangzhou-vpc.dlf.aliyuncs.com

catalog.properties.warehouse

DLF Catalog名稱。

dlf_test

配置樣本

下面為您介紹幾種典型的通過Flink CDC YAML作業將資料同步到資料湖DLF的配置方案:

MySQL整庫同步資料湖DLF

MySQL整庫同步資料到DLF的CDC YAML作業如下所示:

source:
  type: mysql
  name: MySQL Source
  hostname: ${mysql.hostname}
  port: ${mysql.port}
  username: ${mysql.username}
  password: ${mysql.password}
  tables: mysql_test.\.*
  server-id: 8601-8604
  #(可選)同步增量階段新建立的表的資料
  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

sink:
  type: paimon
  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 Source配置】,建議設定下列配置項,詳情請參見MySQL

  1. 參數:scan.binlog.newly-added-table.enabled

    作用:同步增量階段新建立的表的資料。

  2. 參數:include-comments.enabled

    作用:同步表注釋和欄位注釋。

  3. 參數:scan.incremental.snapshot.unbounded-chunk-first.enabled

    作用:避免可能出現的TaskManager OutOfMemory問題。

  4. 參數:scan.only.deserialize.captured.tables.changelog.enabled: true

    作用:僅對作業匹配的表的資料進行解析,加速讀取。

說明

【Paimon Sink配置】

  1. Catalog串連參數

  • 參數首碼:catalog.properties

  • 作用:catalog串連資訊

  1. 建表參數

  • 參數首碼:table.properties

  • 作用:建表資訊

  • 配置建議:在建表參數中添加deletion-vectors.enabled的配置,在不損失太大寫入更新效能的同時,獲得極大的讀取效能提升,達到近即時更新與極速查詢的效果

  • 補充說明:在DLF中已經提供了自動進行檔案合并的功能,不建議在建表參數中添加檔案合并和bucket相關參數,例如bucket、num-sorted-run.compaction-trigger等。

  1. 提交使用者

  • 參數名稱:commit.user

  • 作用:寫入Paimon的檔案提交使用者

  • 配置建議:為不同的作業設定不同的提交使用者,可以設定為作業名

  • 補充說明:預設的提交使用者為admin,在多個作業寫入同一張表時可能出現提交衝突和不一致的問題。

寫入資料湖DLF分區表

資料攝入作業的源表通常不包含分區欄位資訊,如果希望寫入的下遊表為分區表,您需要通過Flink CDC資料攝入作業開發參考中的partition-keys設定分區欄位,配置樣本如下:

source:
  type: mysql
  name: MySQL Source
  hostname: ${mysql.hostname}
  port: ${mysql.port}
  username: ${mysql.username}
  password: ${mysql.password}
  tables: mysql_test.\.*
  server-id: 8601-8604
  #(可選)同步增量階段新建立的表的資料
  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

sink:
  type: paimon
  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

transform:
  - source-table: mysql_test.tbl1
    #(可選)設定分區欄位  
    partition-keys: id,pt
  - source-table: mysql_test.tbl2
    partition-keys: id,pt

寫入資料湖DLF Append Only表

資料攝入作業的源表包含完整的變更類型,如果希望寫入的下遊表為將刪除操作轉化為插入操作實現邏輯刪除的功能,您可以通過Flink CDC資料攝入作業開發參考實現該需求,配置樣本如下:

source:
  type: mysql
  name: MySQL Source
  hostname: ${mysql.hostname}
  port: ${mysql.port}
  username: ${mysql.username}
  password: ${mysql.password}
  tables: mysql_test.\.*
  server-id: 8601-8604
  #(可選)同步增量階段新建立的表的資料
  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

sink:
  type: paimon
  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
  
transform:
  - source-table: mysql_test.tbl1
    #(可選)設定分區欄位
    partition-keys: id,pt
    #(可選)實現虛刪除
    projection: \*, __data_event_type__ AS op_type
    converter-after-transform: SOFT_DELETE
  - source-table: mysql_test.tbl2
    #(可選)設定分區欄位
    partition-keys: id,pt
    #(可選)實現虛刪除
    projection: \*, __data_event_type__ AS op_type
    converter-after-transform: SOFT_DELETE
說明
  • 通過在projection中添加__data_event_type,將變更類型作為新增欄位寫入到下遊表中。同時設定converter-after-transform為SOFT_DELETE,可以將刪除操作轉化為插入操作,使得下遊能夠完整記錄全部變更操作。詳見Flink CDC資料攝入作業開發參考

Kafka CDC資料即時同步到資料湖DLF

假設Kafka的inventory Topic中儲存了兩張表(customers和products)的變更資料,且資料格式為Debezium JSON。以下樣本作業可將這兩張表的資料分別同步到DLF對應的目標表:

source:
  type: kafka
  name: Kafka Source
  properties.bootstrap.servers: ${kafka.bootstrap.servers}
  topic: inventory
  scan.startup.mode: earliest-offset
  value.format: debezium-json
  debezium-json.distributed-tables: true

sink:
  type: paimon
  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

# debezium-json不包含主鍵資訊,需要另外為表添加主鍵  
transform:
  - source-table: \.*.\.*
    projection: \*
    primary-keys: id
說明
  • Kafka資料來源讀取的格式支援canal-json、debezium-json(預設)和json格式。

  • 當資料格式為debezium-json時,由於debezium-json訊息不記錄主鍵資訊,需要通過transform規則手動為表添加主鍵:

    transform:
      - source-table: \.*.\.*
        projection: \*
        primary-keys: id
  • 當單表的資料分布在多個分區中,或資料位元於不同分區中的表需要進行分庫分表合并時,需要將配置項debezium-json.distributed-tablescanal-json.distributed-tables設為true。

  • kafka資料來源支援多種Schema推導策略,可以通過配置項schema.inference.strategy設定,Schema推導和變更同步策略詳情請參見訊息佇列Kafka

Kafka 日誌資料即時同步到資料湖DLF

如果您的Kafka叢集中儲存的是自訂的JSON格式的資料,您可以配置CDC YAML作業同步Kafka資料到DLF儲存,我們會為您提供自動資料類型推導、表結構推導和表結構演化的支援。

假設Kafka的inventory Topic中儲存了一張日誌表的資料,且資料格式為JSON。以下樣本作業可將這張表的資料同步到DLF對應的目標表:

source:
  type: kafka
  name: Kafka Source
  properties.bootstrap.servers: ${kafka.bootstrap.servers}
  topic: inventory
  scan.startup.mode: earliest-offset
  value.format: json
  # (可選)遞迴式地展開JSON中的嵌套列
  json.infer-schema.flatten-nested-columns.enable: true
  # (可選)跳過前 100 次出現的解析異常;若超過 100 次則作業失敗。
  ingestion.ignore-errors: true
  ingestion.error-tolerance.max-count: 100

sink:
  type: paimon
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  #(可選)開啟刪除向量,提升讀取效能
  table.properties.deletion-vectors.enabled: true

# 為表添加主鍵資訊 
transform:
  - source-table: \.*.\.*
    projection: \*
    primary-keys: id
    
# 將 inventory topic 中所有的資料都寫入到 test_database.inventory 表中
route:
  - source-table: inventory
    sink-table: test_database.inventory
    
pipeline:
  # (可選)將會導致處理異常的髒資料記錄到日誌中
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger

假設Kafka的inventory Topic中儲存了多張日誌表的資料,資料格式為JSON,且在JSON內容中的databaseName、tableName欄位中提供了庫名、表名資訊。以下樣本作業可將這個Topic中的多張表的資料同步到DLF對應的目標表:

source:
  type: kafka
  name: Kafka Source
  properties.bootstrap.servers: ${kafka.bootstrap.servers}
  topic: inventory
  scan.startup.mode: earliest-offset
  value.format: json
  # (可選)遞迴式地展開JSON中的嵌套列
  json.infer-schema.flatten-nested-columns.enable: true
  # 使用 databaseName 欄位中的值作為庫名,使用 tableName 欄位中的值作為表名
  json.decode.parser-table-id.fields: databaseName,tableName
  # (可選)跳過前 100 次出現的解析異常;若超過 100 次則作業失敗。
  ingestion.ignore-errors: true
  ingestion.error-tolerance.max-count: 100

sink:
  type: paimon
  catalog.properties.metastore: rest
  catalog.properties.uri: dlf_uri
  catalog.properties.warehouse: your_warehouse
  catalog.properties.token.provider: dlf
  #(可選)開啟刪除向量,提升讀取效能
  table.properties.deletion-vectors.enabled: true

# 為表添加主鍵資訊 
transform:
  - source-table: \.*.\.*
    projection: \*
    primary-keys: id
    
# 將 ods.inventory、ods.customer,ods.user 中的資料分別寫入到 test_database.inventory,test_database.customer,test_database.user 表中
route:
  - source-table: ods.inventory
    sink-table: test_database.inventory
  - source-table: ods.customer
    sink-table: test_database.customer
  - source-table: ods.user
    sink-table: test_database.user   
    
pipeline:
  # (可選)將會導致處理異常的髒資料記錄到日誌中
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger

如果您希望瞭解更多JSON格式的Kafka源表的表結構推導與演化策略,可以參考表結構解析和變更同步策略說明。

如果您希望添加進行更精細的作業配置,可以查詢Flink CDC資料攝入作業開發參考