全部產品
Search
文件中心

Realtime Compute for Apache Flink:日誌資料即時入湖

更新時間:Apr 22, 2026

本文為您介紹使用資料攝入CDC YAML作業將日誌資料寫入阿里雲資料湖構建的最佳實務。

Kafka日誌資料入湖

Apache Kafka是一款開源的分布式訊息佇列系統,廣泛用於高效能資料處理、流式分析、Data Integration等巨量資料領域。通過簡單的YAML作業的編寫,使用者可以使用Flink CDC資料攝入快速完成日誌資料的即時入湖,自動協助使用者完成Schema推導,提供Schema Evolution能力。

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

source:
  type: kafka
  name: Kafka Source
  # Kafka broker地址
  properties.bootstrap.servers: ${kafka.bootstrap.servers}
  # 設定消費的Topic
  topic: inventory
  # 指定從最早的位點開始消費
  scan.startup.mode: earliest-offset
  # 設定Kafka訊息Value部分格式
  value.format: json
  # Schema解析原則設定為continuous,檢測每條訊息Schema並同步Schema變更
  schema.inference.strategy: continuous
  # (可選)遞迴式地展開JSON中的嵌套列
  json.infer-schema.flatten-nested-columns.enable: true
  # (可選)跳過前 100 次出現的解析異常;若超過 100 次則作業失敗。
  ingestion.ignore-errors: true
  ingestion.error-tolerance.max-count: 100

sink:
  type: paimon
  # Metastore類型,固定為rest
  catalog.properties.metastore: rest
  # Token提供方,固定為dlf
  catalog.properties.token.provider: dlf
  # 訪問DLF Rest Catalog Server的URI,格式為http://[region-id]-vpc.dlf.aliyuncs.com,如http://cn-hangzhou-vpc.dlf.aliyuncs.com
  catalog.properties.uri: dlf_uri
  # DLF Catalog名稱。
  catalog.properties.warehouse: your_warehouse
  #(可選)開啟刪除向量,提升讀取效能
  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
  # Metastore類型,固定為rest
  catalog.properties.metastore: rest
  # Token提供方,固定為dlf
  catalog.properties.token.provider: dlf
  # 訪問DLF Rest Catalog Server的URI,格式為http://[region-id]-vpc.dlf.aliyuncs.com,如http://cn-hangzhou-vpc.dlf.aliyuncs.com
  catalog.properties.uri: dlf_uri
  # DLF Catalog名稱。
  catalog.properties.warehouse: your_warehouse
  #(可選)開啟刪除向量,提升讀取效能
  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源表的表結構推導與演化策略,可以參考表結構解析和變更同步策略說明。

典型使用情境

資料攝入Kafka連接器有較多可選擇的配置項,下面將介紹一些典型情境下該如何配置作業,更詳細的配置請參見資料攝入Kafka連接器。

解決欄位名稱衝突

Kafka訊息包含Key和Value兩個部分,使用者可以分別設定key.format和value.format定義兩者的格式,最終Schema由這兩個部分的全部欄位合并組成。

如果Key和Value兩個部分存在欄位名稱衝突,可以通過key.fields-prefix和value.fields-prefix為欄位名稱添加對應首碼解決。

例如Kafka訊息中Key部分包含id和name兩個欄位,Value部分包含id和price兩個欄位,如下作業得到的Schema包含的欄位為key_id,key_name,val_id和val_price。

source:
  type: kafka
  properties.bootstrap.servers: localhost:9092
  topic: test_topic
  properties.group.id: test_group
  scan.startup.mode: earliest-offset
  value.format: json
  key.format: json
  # 為Key中欄位名稱添加首碼
  key.fields-prefix: key_
  # 為Value中欄位名稱添加首碼
  value.fields-prefix: val_
  # Schema解析原則設定為continuous,檢測每條訊息Schema並同步Schema變更
  schema.inference.strategy: continuous
  # 忽略資料解析過程中的報錯
  ingestion.ignore-errors: true
  # 資料解析出現100次報錯後,觸發作業失敗
  ingestion.error-tolerance.max-count: 100

# 將 test_topic topic 中所有的資料都寫入到 test_database.test_topic 表中
route:
  - source-table: test_topic
    sink-table: test_database.test_topic

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
  
pipeline:
  # 開啟髒資料收集,髒資料統一寫入記錄檔
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger

讀取中繼資料

使用metadata.list配置可以讀取額外添加的Kafka訊息中繼資料並傳遞,metadata.list中新增的中繼資料列可以在transform模組中直接使用,更多支援的中繼資料請見可用的中繼資料列。

如下配置會把partition和offset中繼資料添加到資料中,並在transform模組直接使用這兩個中繼資料,過濾出分區大於1且offset大於100的資料。

source:
  type: kafka
  name: Kafka source
  properties.bootstrap.servers: localhost:9092
  topic: test_topic
  properties.group.id: test_group
  scan.startup.mode: earliest-offset
  value.format: json
  # 添加中繼資料列
  metadata.list: partition,offset
  # Schema解析原則設定為continuous,檢測每條訊息Schema並同步Schema變更
  schema.inference.strategy: continuous
  # 忽略資料解析過程中的報錯
  ingestion.ignore-errors: true
  # 資料解析出現100次報錯後,觸發作業失敗
  ingestion.error-tolerance.max-count: 100
  
transform:
  - source-table: \.*.\.*
    filter: '`partition` > 1 and `offset` > 100'
    
# 將 test_topic topic 中所有的資料都寫入到 test_database.test_topic 表中
route:
  - source-table: test_topic
    sink-table: test_database.test_topic

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
  
pipeline:
  # 開啟髒資料收集,髒資料統一寫入記錄檔
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger

髒資料處理

日誌資料可能包含格式不正確的髒資料,這些髒資料會導致作業報錯不斷重啟。Flink CDC資料攝入支援忽略解析報錯並將錯誤資料進行收集,詳情請參見髒資料收集。

如下所示的作業,會將解析錯誤的髒資料統一寫入記錄檔,並在超過100次解析報錯後觸發作業失敗。

source:
  type: kafka
  name: Kafka source
  properties.bootstrap.servers: localhost:9092
  topic: test_topic
  properties.group.id: test_group
  scan.startup.mode: earliest-offset
  value.format: json
  # Schema解析原則設定為continuous,檢測每條訊息Schema並同步Schema變更
  schema.inference.strategy: continuous
  # 忽略資料解析過程中的報錯
  ingestion.ignore-errors: true
  # 資料解析出現100次報錯後,觸發作業失敗
  ingestion.error-tolerance.max-count: 100
  
# 將 test_topic topic 中所有的資料都寫入到 test_database.test_topic 表中
route:
  - source-table: test_topic
    sink-table: test_database.test_topic

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
  
pipeline:
  # 開啟髒資料收集,髒資料統一寫入記錄檔
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger

如果完全不需要收集髒資料資訊,並且不希望作業因髒資料而發生失敗,可以使用json.ignore-parse-errors,debezium-json.ignore-parse-errors或canal-json.ignore-parse-errors直接忽略解析報錯。

指定TableID

Debezium JSON和Canal JSON格式是固定的,格式中有固定的欄位儲存TableID內容。JSON資料欄位無固定格式,預設使用topic名稱作為TableID,如果需要指定欄位中的資料作為TableID,JSON格式需要配置json.decode.parser-table-id.fields,例如:JSON資料為{"col0":"a", "col1":"b", "col2":"c"},不同配置產生TableID如下:

配置

tableId

col0

a

col0,col1

a.b

col0,col1,col2

a.b.c

如下作業配置會將資料中的db和table列拼接作為最終的TableID。

source:
  type: kafka
  name: Kafka source
  properties.bootstrap.servers: localhost:9092
  topic: test_topic
  properties.group.id: test_group
  scan.startup.mode: earliest-offset
  value.format: json
  # 使用資料中的db和table欄位作為table id
  json.decode.parser-table-id.fields: db,table
  # Schema解析原則設定為continuous,檢測每條訊息Schema並同步Schema變更
  schema.inference.strategy: continuous
  # 忽略資料解析過程中的報錯
  ingestion.ignore-errors: true
  # 資料解析出現100次報錯後,觸發作業失敗
  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
  
pipeline:
  # 開啟髒資料收集,髒資料統一寫入記錄檔
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger

資料類型推導

日誌資料的類型需要通過解析資料來進行推斷,詳細的類型推導請參見表結構解析和變更同步策略。下面介紹一些常見情境下,如何調整配置控制類型解析。

(通用)欄位類型統一設定成String

如果下遊處理和儲存不需要關心欄位類型,可以開啟json.infer-schema.primitive-as-string,debezium-json.infer-schema.primitive-as-string或canal-json.infer-schema.primitive-as-string將欄位類型推導跳過並統一設定為String類型。

source:
  type: kafka
  name: Kafka source
  properties.bootstrap.servers: localhost:9092
  topic: test_topic
  properties.group.id: test_group
  scan.startup.mode: earliest-offset
  value.format: json
  # JSON資料中所有欄位類型使用String
  json.infer-schema.primitive-as-string: true
  # Schema解析原則設定為continuous,檢測每條訊息Schema並同步Schema變更
  schema.inference.strategy: continuous
  # 忽略資料解析過程中的報錯
  ingestion.ignore-errors: true
  # 資料解析出現100次報錯後,觸發作業失敗
  ingestion.error-tolerance.max-count: 100

# 將 test_topic topic 中所有的資料都寫入到 test_database.test_topic 表中
route:
  - source-table: test_topic
    sink-table: test_database.test_topic

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
  
pipeline:
  # 開啟髒資料收集,髒資料統一寫入記錄檔
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger

(通用)指定初始Schema

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

source:
  type: kafka
  name: Kafka source
  properties.bootstrap.servers: localhost:9092
  topic: test_topic
  properties.group.id: test_group
  scan.startup.mode: earliest-offset
  value.format: json
  # 使用資料中的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);
  # Schema解析原則設定為continuous,檢測每條訊息Schema並同步Schema變更
  schema.inference.strategy: continuous
  # 忽略資料解析過程中的報錯
  ingestion.ignore-errors: true
  # 資料解析出現100次報錯後,觸發作業失敗
  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
  
pipeline:
  # 開啟髒資料收集,髒資料統一寫入記錄檔
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger

如上配置為 db1.t1 表指定 id 欄位的初始類型為 BIGINT,name 欄位的初始類型為 VARCHAR(10) ,為 db1.t2 表指定 id 欄位的初始類型為 BIGINT。

(JSON)固定欄位類型

JSON資料的欄位類型需要通過解析JSON資料節點類型進行推斷,有時推斷出的類型並不符合使用者期望,此時可以通過配置json.infer-schema.fixed-types為某些欄位指定類型。

如下的作業配置指定id欄位類型為BIGINT,name欄位類型為VARCHAR(10)。

source:
  type: kafka
  name: Kafka source
  properties.bootstrap.servers: localhost:9092
  topic: test_topic
  properties.group.id: test_group
  scan.startup.mode: earliest-offset
  value.format: json
  # 設定特定欄位始終為固定類型
  json.infer-schema.fixed-types: id BIGINT, name VARCHAR(10)
  # 11.5版本及以前版本需要搭配此配置使用
  scan.max.pre.fetch.records: 0
  # Schema解析原則設定為continuous,檢測每條訊息Schema並同步Schema變更
  schema.inference.strategy: continuous
  # 忽略資料解析過程中的報錯
  ingestion.ignore-errors: true
  # 資料解析出現100次報錯後,觸發作業失敗
  ingestion.error-tolerance.max-count: 100

# 將 test_topic topic 中所有的資料都寫入到 test_database.test_topic 表中
route:
  - source-table: test_topic
    sink-table: test_database.test_topic
    
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
  
pipeline:
  # 開啟髒資料收集,髒資料統一寫入記錄檔
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger

(Canal JSON)指定推斷來源

Canal JSON資料相比JSON資料記錄了更多的資訊,如果Canal JSON資料中sqlType或mysqlType欄位存在,可以通過這兩部分的資訊解析出更加精確的資料類型。

  • 由sqlType解析Schema

source:
  type: kafka
  name: Kafka source
  properties.bootstrap.servers: localhost:9092
  topic: test_topic
  properties.group.id: test_group
  scan.startup.mode: earliest-offset
  value.format: canal-json
  # 通過sqlType欄位推斷Schema
  canal-json.infer-schema.strategy: SQL_TYPE
  # Schema解析原則設定為continuous,檢測每條訊息Schema並同步Schema變更
  schema.inference.strategy: continuous
  # 忽略資料解析過程中的報錯
  ingestion.ignore-errors: true
  # 資料解析出現100次報錯後,觸發作業失敗
  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
  
pipeline:
  # 開啟髒資料收集,髒資料統一寫入記錄檔
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger
  • 由mysqlType解析Schema

    source:
      type: kafka
      name: Kafka source
      properties.bootstrap.servers: localhost:9092
      topic: test_topic
      properties.group.id: test_group
      scan.startup.mode: earliest-offset
      value.format: canal-json
      # 通過mysqlType欄位推斷Schema
      canal-json.infer-schema.strategy: MYSQL_TYPE
      # Schema解析原則設定為continuous,檢測每條訊息Schema並同步Schema變更
      schema.inference.strategy: continuous
      # 忽略資料解析過程中的報錯
      ingestion.ignore-errors: true
      # 資料解析出現100次報錯後,觸發作業失敗
      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
      
    pipeline:
      # 開啟髒資料收集,髒資料統一寫入記錄檔
      dirty-data.collector:
        name: Logger Dirty Data Collector
        type: logger

使用靜態Schema

如果Topic中資料的Schema是固定不變的,可以將schema.inference.strategy設定為static,這樣配置後資料攝入作業僅會在作業啟動時進行一次Schema推導,並且對後續資料也不會再進行Schema解析。

source:
  type: kafka
  properties.bootstrap.servers: localhost:9092
  topic: test_topic
  properties.group.id: test_group
  scan.startup.mode: earliest-offset
  value.format: json
  # Schema解析原則設定為static,僅在作業啟動時推導一次Schema
  schema.inference.strategy: static
  # 嘗試在每個分區消費20條資料來推導Schema
  scan.max.pre.fetch.records: 20
  # 忽略資料解析過程中的報錯
  ingestion.ignore-errors: true
  # 資料解析出現100次報錯後,觸發作業失敗
  ingestion.error-tolerance.max-count: 100
  
# 將 test_topic topic 中所有的資料都寫入到 test_database.test_topic 表中
route:
  - source-table: test_topic
    sink-table: test_database.test_topic
    
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
  
pipeline:
  # 開啟髒資料收集,髒資料統一寫入記錄檔
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger

此外還可以通過添加scan.value.initial-schemas.ddls參數指定初始表結構,跳過某些表的Schema推導,如下為表db1.t1和db1.t2指定了初始表結構。

source:
  type: kafka
  name: Kafka source
  properties.bootstrap.servers: localhost:9092
  topic: test_topic
  properties.group.id: test_group
  scan.startup.mode: earliest-offset
  value.format: json
  # 使用資料中的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);
  # Schema解析原則設定為static,僅在作業啟動時推導一次Schema
  schema.inference.strategy: static
  # 忽略資料解析過程中的報錯
  ingestion.ignore-errors: true
  # 資料解析出現100次報錯後,觸發作業失敗
  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
  
pipeline:
  # 開啟髒資料收集,髒資料統一寫入記錄檔
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger

Kafka日誌資料解析加速

除了通用的Kafka連接器加速配置,針對使用情境的不同,Flink CDC資料攝入作業有一些可以加速解析的配置項。

  1. 解析Schema是耗時較多的操作,如果下遊只需要String類型,可以開啟json.infer-schema.primitive-as-string,debezium-json.infer-schema.primitive-as-string或canal-json.infer-schema.primitive-as-string將欄位類型推導跳過並統一設定為String類型來加速解析。

  2. Canal JSON資料可以使用canal-json.database.include和canal-json.table.include來過濾掉不需要的表的資料。

  3. 如果可以確定資料Schema不會發生變化,或者不需要進行Schema Evolution,可以將schema.inference.strategy修改為static,只在第一次作業啟動時進行Schema推斷。

SLS日誌資料入湖

SLSLog Service是針對日誌類資料的一站式服務,Flink CDC資料攝入支援快速完成SLS日誌資料的即時入湖,自動協助使用者完成Schema推導,提供Schema Evolution能力。

source:
  type: sls
  endpoint: localhost
  project: test_pj
  logstore: test_log
  accessId: access_id
  accessKey: access_key
  
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
  
pipeline:
  # 開啟髒資料收集,髒資料統一寫入記錄檔
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger
  • TableID:預設會拼接project和logstore作為table id,如上作業得到的table id為test_pj.test_log。

  • 資料類型:SLS連接器預設會將每條日誌中的欄位類型都推導為String類型。

  • Schema Evolution:目前僅支援新增列變更,預設會在Schema結尾添加新列。

典型使用情境

下面將介紹一些典型情境下該如何配置作業,更詳細的配置請參見資料攝入SLS連接器。

讀取中繼資料

使用metadata.list配置可以讀取額外添加的SLS中繼資料並傳遞,metadata.list中新增的中繼資料列可以在transform模組中直接使用,更多支援的中繼資料請見metadata.list。

需要注意metadata.list添加的列並不會直接放入資料的列中,如果將這部分中繼資料列寫出,需要在transform的projection部分聲明出這些列。

如下配置會把__timestamp__和__tag__中繼資料添加到資料並寫入下遊,並在transform模組使用這兩個中繼資料,過濾出時間戳記大於1772181154且tag為test的資料。

source:
  type: sls
  endpoint: localhost
  project: test_pj
  logstore: test_log
  accessId: access_id
  accessKey: access_key
  # 添加中繼資料列
  metadata.list: __timestamp__,__tag__
  # Schema解析原則設定為continuous,檢測每條訊息Schema並同步Schema變更
  schema.inference.strategy: continuous
  # 忽略資料解析過程中的報錯
  ingestion.ignore-errors: true
  # 資料解析出現100次報錯後,觸發作業失敗
  ingestion.error-tolerance.max-count: 100
  
transform:
  - source-table: \.*.\.*
    projection: \*, __timestamp__ as timestamp_col, __tag__ as tag_col
    filter: '`__timestamp__` > 1772181154 and `__tag__` = "test"'

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
  
pipeline:
  # 開啟髒資料收集,髒資料統一寫入記錄檔
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger

髒資料處理

日誌資料可能包含格式不正確的髒資料,這些髒資料會導致作業報錯不斷重啟。Flink CDC資料攝入支援忽略解析報錯並將錯誤資料進行收集,詳情請參見髒資料收集。

如下所示的作業,會將解析錯誤的髒資料統一寫入記錄檔,並在超過100次解析報錯後觸發作業失敗。

source:
  type: sls
  endpoint: localhost
  project: test_pj
  logstore: test_log
  accessId: access_id
  accessKey: access_key
  # Schema解析原則設定為continuous,檢測每條訊息Schema並同步Schema變更
  schema.inference.strategy: continuous
  # 忽略資料解析過程中的報錯
  ingestion.ignore-errors: true
  # 資料解析出現100次報錯後,觸發作業失敗
  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
  
pipeline:
  # 開啟髒資料收集,髒資料統一寫入記錄檔
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger

指定TableID

如果需要指定欄位中的資料作為TableID,需要配置decode.table-id.fields,例如:日誌資料為{"col0":"a", "col1":"b", "col2":"c"},不同配置產生TableID如下:

配置

tableId

col0

a

col0,col1

a.b

col0,col1,col2

a.b.c

如下作業配置會將資料中的db和table列拼接作為最終的TableID。

source:
  type: sls
  endpoint: localhost
  project: test_pj
  logstore: test_log
  accessId: access_id
  accessKey: access_key
  # 使用資料中的db和table欄位作為table id
  decode.table-id.fields: db,table
  # Schema解析原則設定為continuous,檢測每條訊息Schema並同步Schema變更
  schema.inference.strategy: continuous
  # 忽略資料解析過程中的報錯
  ingestion.ignore-errors: true
  # 資料解析出現100次報錯後,觸發作業失敗
  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
  
pipeline:
  # 開啟髒資料收集,髒資料統一寫入記錄檔
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger

指定欄位類型

資料攝入SLS連接器預設將資料中的全部欄位看作String類型, 如果需要指定其中某些欄位的類型,可以通過fixed-types配置項指定。

如下作業可以指定id欄位類型為BIGINT,name欄位類型為VARCHAR(10)。

source:
  type: sls
  endpoint: localhost
  project: test_pj
  logstore: test_log
  accessId: access_id
  accessKey: access_key
  # 指定id欄位類型為BIGINT,name欄位類型為VARCHAR(10)
  fixed-types: id BIGINT, name VARCHAR(10)
  # Schema解析原則設定為continuous,檢測每條訊息Schema並同步Schema變更
  schema.inference.strategy: continuous
  # 忽略資料解析過程中的報錯
  ingestion.ignore-errors: true
  # 資料解析出現100次報錯後,觸發作業失敗
  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
  
pipeline:
  # 開啟髒資料收集,髒資料統一寫入記錄檔
  dirty-data.collector:
    name: Logger Dirty Data Collector
    type: logger