如何在阿里雲Realtime ComputeFlink版上通過Flink CDC以Paimon REST訪問DLF Catalog。
前提條件
已建立Flink全託管工作空間。如未建立,詳情請參見開通Realtime ComputeFlink版。
請確保Flink工作空間與DLF位於同一個VPC下。
使用限制
僅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配置項說明如下:
配置項 | 描述 | 是否必填 | 樣本 |
| Metastore類型,固定為rest。 | 是 | rest |
| Token提供方,固定為dlf。 | 是 | dlf |
| 訪問DLF Rest Catalog Server的URI,格式為 | 是 | http://cn-hangzhou-vpc.dlf.aliyuncs.com |
| 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。
參數:scan.binlog.newly-added-table.enabled
作用:同步增量階段新建立的表的資料。
參數:include-comments.enabled
作用:同步表注釋和欄位注釋。
參數:scan.incremental.snapshot.unbounded-chunk-first.enabled
作用:避免可能出現的TaskManager OutOfMemory問題。
參數:scan.only.deserialize.captured.tables.changelog.enabled: true
作用:僅對作業匹配的表的資料進行解析,加速讀取。
【Paimon Sink配置】
Catalog串連參數
參數首碼:catalog.properties
作用:catalog串連資訊
建表參數
參數首碼:table.properties
作用:建表資訊
配置建議:在建表參數中添加deletion-vectors.enabled的配置,在不損失太大寫入更新效能的同時,獲得極大的讀取效能提升,達到近即時更新與極速查詢的效果
補充說明:在DLF中已經提供了自動進行檔案合并的功能,不建議在建表參數中添加檔案合并和bucket相關參數,例如bucket、num-sorted-run.compaction-trigger等。
提交使用者
參數名稱: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即時同步到資料湖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: idKafka資料來源讀取的格式支援canal-json、debezium-json(預設)和json格式。
當資料格式為debezium-json時,由於debezium-json訊息不記錄主鍵資訊,需要通過transform規則手動為表添加主鍵:
transform: - source-table: \.*.\.* projection: \* primary-keys: id當單表的資料分布在多個分區中,或資料位元於不同分區中的表需要進行分庫分表合并時,需要將配置項debezium-json.distributed-tables或canal-json.distributed-tables設為true。
kafka資料來源支援多種Schema推導策略,可以通過配置項schema.inference.strategy設定,Schema推導和變更同步策略詳情請參見訊息佇列Kafka。
如果您希望添加進行更精細的作業配置,可以查詢Flink CDC資料攝入作業開發參考。