本文介紹如何使用Iceberg連接器。
背景資訊
Apache Iceberg是一種開放的資料湖表格格式。您可以藉助Apache Iceberg快速地在HDFS或者雲端OSS上構建自己的資料湖儲存服務,並藉助開源巨量資料生態的Flink、Spark、Hive、Presto等計算引擎來實現資料湖的分析。
|
類別 |
詳情 |
|
支援類型 |
源表和結果表,資料攝入目標端 |
|
運行模式 |
批模式和流模式 |
|
資料格式 |
暫不適用 |
|
特有監控指標 |
暫無 |
|
API種類 |
SQL,資料攝入YAML作業 |
|
是否支援更新或刪除結果表資料 |
是 |
特色功能
目前Apache Iceberg提供以下核心能力:
-
基於HDFS或者Object Storage Service構建低成本的輕量級資料湖儲存服務。
-
完善的ACID語義。
-
支援歷史版本回溯。
-
支援高效的資料過濾。
-
支援Schema Evolution。
-
支援Partition Evolution。
您可以藉助Flink高效的容錯能力和流處理能力,把海量的日誌行為資料即時匯入到Apache Iceberg資料湖內,再藉助Flink或者其他分析引擎來實現資料價值的提取。
使用限制
-
僅Flink計算引擎VVR 4.0.8及以上版本支援Iceberg連接器。Iceberg連接器需要搭配DLF Catalog一起使用,詳情請參見管理DLF-Legacy Catalog。
-
Iceberg連接器支援Apache Iceberg v1和v2表格式,詳情請參見Iceberg Table Spec。
說明僅Realtime Compute引擎VVR 8.0.7及以上版本支援v2表格式。
-
流讀模式下,僅支援將Append Only的Iceberg表作為源表。
文法結構
CREATE TABLE iceberg_table (
id BIGINT,
data STRING
PRIMARY KEY(`id`) NOT ENFORCED
)
PARTITIONED BY (data)
WITH (
'connector' = 'iceberg',
...
);
WITH參數
通用(源表)
|
參數 |
說明 |
資料類型 |
是否必填 |
預設值 |
備忘 |
|
connector |
源表類型 |
String |
是 |
無 |
固定值為 |
|
catalog-name |
Catalog名稱 |
String |
是 |
無 |
請填寫為自訂的英文名。 |
|
catalog-database |
資料庫名稱 |
String |
是 |
default |
對應在DLF上建立的資料庫名稱,例如dlf_db。 說明
如果您沒有建立對應的DLF資料庫,請建立DLF資料庫。 |
|
io-impl |
Distributed File System的實作類別名 |
String |
是 |
無 |
固定值為 |
|
oss.endpoint |
阿里雲Object Storage Service服務OSS的Endpoint |
String |
否 |
無 |
請詳情參見地區和Endpoint。 說明
|
|
阿里雲帳號的AccessKey ID |
String |
是 |
無 |
詳情請參見如何查看AccessKey ID和AccessKey Secret資訊? 重要
為了避免您的AK資訊泄露,建議您使用變數的方式填寫AccessKey取值,詳情請參見專案變數。 |
|
阿里雲帳號的AccessKey Secret |
String |
是 |
無 |
|
|
catalog-impl |
Catalog的Class類名 |
String |
是 |
無 |
固定值為 |
|
warehouse |
表資料存放在OSS的路徑 |
String |
是 |
無 |
無。 |
|
dlf.catalog-id |
阿里雲帳號的帳號ID |
String |
是 |
無 |
可通過使用者資訊頁面擷取帳號ID。 |
|
dlf.endpoint |
DLF服務的Endpoint |
String |
是 |
無 |
。 說明
|
|
dlf.region-id |
DLF服務的地區名 |
String |
是 |
無 |
。 說明
請和dlf.endpoint選擇的地區保持一致。 |
結果表專屬
|
參數 |
說明 |
資料類型 |
是否必填 |
預設值 |
備忘 |
|
write.operation |
寫入操作模式 |
String |
否 |
upsert |
|
|
hive_sync.enable |
是否開啟同步中繼資料到Hive功能 |
boolean |
否 |
false |
參數取值如下:
|
|
hive_sync.mode |
Hive資料同步模式 |
String |
否 |
hms |
|
|
hive_sync.db |
同步到Hive的資料庫名稱 |
String |
否 |
當前Table在Catalog中的資料庫名 |
無。 |
|
hive_sync.table |
同步到Hive的表名稱 |
String |
否 |
當前Table名 |
無。 |
|
dlf.catalog.region |
DLF服務的地區名 |
String |
否 |
無 |
。 說明
|
|
dlf.catalog.endpoint |
DLF服務的Endpoint |
String |
否 |
無 |
。 說明
|
類型映射
|
Iceberg欄位類型 |
Flink欄位類型 |
|
BOOLEAN |
BOOLEAN |
|
INT |
INT |
|
LONG |
BIGINT |
|
FLOAT |
FLOAT |
|
DOUBLE |
DOUBLE |
|
DECIMAL(P,S) |
DECIMAL(P,S) |
|
DATE |
DATE |
|
TIME |
TIME 說明
Iceberg時間戳記精度為微秒,Flink時間戳記精度為毫秒。在使用Flink讀取Iceberg資料時,時間精度會對齊到毫秒。 |
|
TIMESTAMP |
TIMESTAMP |
|
TIMESTAMPTZ |
TIMESTAMP_LTZ |
|
STRING |
STRING |
|
FIXED(L) |
BYTES |
|
BINARY |
VARBINARY |
|
STRUCT<...> |
ROW |
|
LIST<E> |
LIST |
|
MAP<K,V> |
MAP |
程式碼範例
請確認您已建立了OSS Bucket和DLF資料庫。詳情請參見控制台建立儲存空間和資料庫表及函數。
在建立DLF資料庫選擇路徑時,建議按照${warehouse}/${database_name}.db格式填寫。例如,如果warehouse地址為oss://iceberg-test/warehouse,資料庫的名稱為dlf_db,則dlf_db的OSS路徑需要設定為oss://iceberg-test/warehouse/dlf_db.db。
結果表示例
本樣本為您介紹如何通過Datagen連接器隨機產生流式資料寫入Iceberg表。
CREATE TEMPORARY TABLE datagen(
id BIGINT,
data STRING
) WITH (
'connector' = 'datagen'
);
CREATE TEMPORARY TABLE dlf_iceberg (
id BIGINT,
data STRING
) WITH (
'connector' = 'iceberg',
'catalog-name' = '<yourCatalogName>',
'catalog-database' = '<yourDatabaseName>',
'io-impl' = 'org.apache.iceberg.aliyun.oss.OSSFileIO',
'oss.endpoint' = '<yourOSSEndpoint>',
'access.key.id' = '${secret_values.ak_id}',
'access.key.secret' = '${secret_values.ak_secret}',
'catalog-impl' = 'org.apache.iceberg.aliyun.dlf.DlfCatalog',
'warehouse' = '<yourOSSWarehousePath>',
'dlf.catalog-id' = '<yourCatalogId>',
'dlf.endpoint' = '<yourDLFEndpoint>',
'dlf.region-id' = '<yourDLFRegionId>'
);
INSERT INTO dlf_iceberg SELECT * FROM datagen;
源表示例
-
使用DLF Catalog,將Iceberg源表資料寫入到Iceberg結果表中。
CREATE TEMPORARY TABLE src_iceberg ( id BIGINT, data STRING ) WITH ( 'connector' = 'iceberg', 'catalog-name' = '<yourCatalogName>', 'catalog-database' = '<yourDatabaseName>', 'io-impl' = 'org.apache.iceberg.aliyun.oss.OSSFileIO', 'oss.endpoint' = '<yourOSSEndpoint>', 'access.key.id' = '${secret_values.ak_id}', 'access.key.secret' = '${secret_values.ak_secret}', 'catalog-impl' = 'org.apache.iceberg.aliyun.dlf.DlfCatalog', 'warehouse' = '<yourOSSWarehousePath>', 'dlf.catalog-id' = '<yourCatalogId>', 'dlf.endpoint' = '<yourDLFEndpoint>', 'dlf.region-id' = '<yourDLFRegionId>' ); CREATE TEMPORARY TABLE dst_iceberg ( id BIGINT, data STRING ) WITH ( 'connector' = 'iceberg', 'catalog-name' = '<yourCatalogName>', 'catalog-database' = '<yourDatabaseName>', 'io-impl' = 'org.apache.iceberg.aliyun.oss.OSSFileIO', 'oss.endpoint' = '<yourOSSEndpoint>', 'access.key.id' = '${secret_values.ak_id}', 'access.key.secret' = '${secret_values.ak_secret}', 'catalog-impl' = 'org.apache.iceberg.aliyun.dlf.DlfCatalog', 'warehouse' = '<yourOSSWarehousePath>', 'dlf.catalog-id' = '<yourCatalogId>', 'dlf.endpoint' = '<yourDLFEndpoint>', 'dlf.region-id' = '<yourDLFRegionId>' ); BEGIN STATEMENT SET; INSERT INTO src_iceberg VALUES (1, 'AAA'), (2, 'BBB'), (3, 'CCC'), (4, 'DDD'), (5, 'EEE'); INSERT INTO dst_iceberg SELECT * FROM src_iceberg; END;
資料攝入
Iceberg連接器可以用於資料攝入YAML作業開發,作為目標端寫入。
文法結構
sink:
type: iceberg
name: Iceberg Sink
catalog.properties.rest.signing-region: cn-beijing
catalog.properties.uri: http://cn-beijing-vpc.dlf.aliyuncs.com/iceberg
catalog.properties.warehouse: flink_iceberg
catalog.properties.type: rest
catalog.properties.io-impl: org.apache.iceberg.rest.DlfFileIO
配置項
|
參數 |
說明 |
是否必填 |
資料類型 |
預設值 |
備忘 |
|
type |
連接器類型。 |
是 |
STRING |
無 |
固定值為 |
|
name |
目標端名稱。 |
否 |
STRING |
無 |
Sink的名稱。 |
|
catalog.properties.rest.signing-region |
DLF的Region ID,詳見服務存取點。 |
是 |
STRING |
無 |
無 |
|
catalog.properties.uri |
訪問DLF Rest Catalog的URI,詳見Iceberg REST。 |
是 |
STRING |
無 |
無 |
|
catalog.properties.warehouse |
DLF Catalog名稱。 |
是 |
STRING |
無 |
無 |
|
catalog.properties.warehouse |
檔案儲存體的根目錄。 |
否 |
STRING |
無 |
無 |
|
catalog.properties.type |
Catalog類型,固定為rest。 |
是 |
STRING |
rest |
無 |
|
catalog.properties.io-impl |
固定值:org.apache.iceberg.rest.DlfFileIO。 |
是 |
STRING |
org.apache.iceberg.rest.DlfFileIO |
無 |
|
partition.key |
每個分區表的分區欄位。 |
否 |
STRING |
無 |
每個分區表的分區鍵,允許為多個表設定多個主鍵。表之間用 對於需要進行隱式轉換的分區,我們可以直接在分區欄位上添加隱式轉換的函數,例如 |
|
table.properties.* |
建立Iceberg table的參數。 |
否 |
String |
無 |
詳情請參見Iceberg table options。 |
複用已有 Catalog
自VVR 11.5版本起,您可以在Flink CDC資料攝入作業中直接引用“資料管理”頁面中建立的內建Iceberg Catalog,減少手寫串連屬性工作量。
sink:
type: iceberg
using.built-in-catalog: iceberg_catalog
目前,資料攝入作業支援自動複用所有 Iceberg Catalog 參數,等價於在 YAML 作業中手動設定 catalog.properties.首碼的參數。
如果希望覆蓋以上自動複用的參數,可顯式寫出相應的 YAML 參數,其具備更高的優先順序。
使用樣本
Iceberg Catalog為DLF Catalog,寫入阿里雲資料湖構建的配置樣本:
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: iceberg name: Iceberg Sink catalog.properties.rest.signing-region: cn-beijing catalog.properties.uri: http://cn-beijing-vpc.dlf.aliyuncs.com/iceberg catalog.properties.warehouse: flink_iceberg catalog.properties.type: rest catalog.properties.io-impl: org.apache.iceberg.rest.DlfFileIO其中,catalog.properties首碼的參數含義請參見建立Iceberg DLF Catalog。
表結構變更
目前,Iceberg作為資料攝入目標端支援以下表結構變更事件:
-
CREATE TABLE EVENT
-
ADD COLUMN EVENT
-
ALTER COLUMN TYPE EVENT(不支援修改主鍵列的類型)
-
RENAME COLUMN EVENT
-
DROP COLUMN EVENT
-
TRUNCATE TABLE EVENT
-
DROP TABLE EVENT
在下遊 Iceberg 表已經存在時,會優先使用已有的表結構進行寫入,不會嘗試重複建表。
相關文檔
Flink支援的連接器,請參見支援的連接器。