本文為您介紹使用資料攝入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: loggerKafka日誌資料解析加速
除了通用的Kafka連接器加速配置,針對使用情境的不同,Flink CDC資料攝入作業有一些可以加速解析的配置項。
解析Schema是耗時較多的操作,如果下遊只需要String類型,可以開啟json.infer-schema.primitive-as-string,debezium-json.infer-schema.primitive-as-string或canal-json.infer-schema.primitive-as-string將欄位類型推導跳過並統一設定為String類型來加速解析。
Canal JSON資料可以使用canal-json.database.include和canal-json.table.include來過濾掉不需要的表的資料。
如果可以確定資料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: loggerTableID:預設會拼接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