全部產品
Search
文件中心

Realtime Compute for Apache Flink:Iceberg

更新時間:Jun 03, 2026

本文介紹如何使用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

固定值為iceberg

catalog-name

Catalog名稱

String

請填寫為自訂的英文名。

catalog-database

資料庫名稱

String

default

對應在DLF上建立的資料庫名稱,例如dlf_db。

說明

如果您沒有建立對應的DLF資料庫,請建立DLF資料庫。

io-impl

Distributed File System的實作類別名

String

固定值為org.apache.iceberg.aliyun.oss.OSSFileIO

oss.endpoint

阿里雲Object Storage Service服務OSS的Endpoint

String

請詳情參見地區和Endpoint

說明
  • 推薦您為oss.endpoint參數配置OSS的VPC Endpoint。例如,如果您選擇的地區為cn-hangzhou地區,則oss.endpoint需要配置為oss-cn-hangzhou-internal.aliyuncs.com。

  • 如果您需要跨VPC訪問OSS,則請參見如何訪問跨VPC的其他服務?

  • access.key.id:VVR 8.0.6及以下版本

  • access-key-id:VVR 8.0.7及以上版本

阿里雲帳號的AccessKey ID

String

詳情請參見如何查看AccessKey ID和AccessKey Secret資訊?

重要

為了避免您的AK資訊泄露,建議您使用變數的方式填寫AccessKey取值,詳情請參見專案變數

  • access.key.secret:VVR 8.0.6及以下版本

  • access-key-secret:VVR 8.0.7及以上版本

阿里雲帳號的AccessKey Secret

String

catalog-impl

Catalog的Class類名

String

固定值為org.apache.iceberg.aliyun.dlf.DlfCatalog

warehouse

表資料存放在OSS的路徑

String

無。

dlf.catalog-id

阿里雲帳號的帳號ID

String

可通過使用者資訊頁面擷取帳號ID。

dlf.endpoint

DLF服務的Endpoint

String

說明
  • 推薦您為dlf.endpoint參數配置DLF的VPC Endpoint。例如,如果您選擇的地區為cn-hangzhou地區,則dlf.endpoint參數需要配置為dlf-vpc.cn-hangzhou.aliyuncs.com

  • 如果您需要跨VPC訪問DLF,則請參見空間管理與操作

dlf.region-id

DLF服務的地區名

String

說明

請和dlf.endpoint選擇的地區保持一致。

結果表專屬

參數

說明

資料類型

是否必填

預設值

備忘

write.operation

寫入操作模式

String

upsert

  • upsert(預設):資料更新。

  • insert:資料追加寫入。

  • bulk_insert:批量寫入(不更新)。

hive_sync.enable

是否開啟同步中繼資料到Hive功能

boolean

false

參數取值如下:

  • true:開啟

  • false(預設值):不開啟。

hive_sync.mode

Hive資料同步模式

String

hms

  • hms(預設值):採用DLF Catalog時,需要設定hms。

  • jdbc:採用jdbc Catalog時,需要設定為jdbc。

hive_sync.db

同步到Hive的資料庫名稱

String

當前Table在Catalog中的資料庫名

無。

hive_sync.table

同步到Hive的表名稱

String

當前Table名

無。

dlf.catalog.region

DLF服務的地區名

String

說明
  • 僅當hive_sync.mode設定為hms時,dlf.catalog.region參數設定才生效。

  • 請和dlf.catalog.endpoint選擇的地區保持一致。

dlf.catalog.endpoint

DLF服務的Endpoint

String

說明
  • 僅當hive_sync.mode設定為hms時,dlf.catalog.endpoint參數設定才生效。

  • 推薦您為dlf.catalog.endpoint參數配置DLF的VPC Endpoint。例如,如果您選擇的地區為cn-hangzhou地區,則dlf.catalog.endpoint參數需要配置為dlf-vpc.cn-hangzhou.aliyuncs.com

  • 如果您需要跨VPC訪問DLF,則請參見空間管理與操作

類型映射

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

固定值為iceberg

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

每個分區表的分區鍵,允許為多個表設定多個主鍵。表之間用;分隔,分區鍵之間用,分隔。例如,我們可以通過 testdb.table1:id1,id2;testdb.table2:name 來設定testdb.table1表的分區欄位為id1testdb.table2表的分區欄位為name

對於需要進行隱式轉換的分區,我們可以直接在分區欄位上添加隱式轉換的函數,例如 testdb.table1:truncate[10](id);testdb.table2:hour(create_time);testdb.table3:day(create_time);testdb.table4:month(create_time);testdb.table5:year(create_time);testdb.table6:bucket[10](create_time)

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支援的連接器,請參見支援的連接器