本文介紹如何通過開源 Apache Iceberg Java API、Spark 和 Flink 訪問 DLF Iceberg REST Catalog,管理表中繼資料並讀寫 Iceberg 表。DLF Iceberg Catalog 相容標準 Iceberg REST 協議,使用 AWS SigV4 簽名認證,資料面通過社區標準 S3FileIO 直連 OSS S3 相容端點。
前提條件
運行環境需使用 JDK 17 及以上版本(Iceberg 1.11.0 要求)。
通過 Maven 引入以下 Apache Iceberg 社區開源依賴(版本 1.11.0),所有依賴均可從 Maven Central 直接拉取:
<dependencies>
<dependency>
<groupId>org.apache.iceberg</groupId>
<artifactId>iceberg-core</artifactId>
<version>1.11.0</version>
</dependency>
<dependency>
<groupId>org.apache.iceberg</groupId>
<artifactId>iceberg-aws</artifactId>
<version>1.11.0</version>
</dependency>
<dependency>
<groupId>org.apache.iceberg</groupId>
<artifactId>iceberg-aws-bundle</artifactId>
<version>1.11.0</version>
</dependency>
</dependencies>如需在本地通過 Java API 直接讀寫 Parquet 資料(而非通過計算引擎),需額外引入 iceberg-data、iceberg-parquet、iceberg-orc、parquet-hadoop 及 hadoop-common 依賴。
使用 Java API 串連
通過 Iceberg 社區標準 RESTCatalog 類串連 DLF Iceberg REST Catalog,初始化時傳入以下配置參數:
Map<String, String> props = new HashMap<>();
props.put("uri", "http://cn-hangzhou-vpc.dlf.aliyuncs.com/iceberg");
props.put("warehouse", "iceberg_table_test");
props.put("rest.auth.type", "sigv4");
props.put("rest.auth.sigv4.delegate-auth-type", "none");
props.put("rest.signing-region", "cn-hangzhou");
props.put("rest.signing-name", "DlfNext");
props.put("rest.access-key-id", "xxx");
props.put("rest.secret-access-key", "yyy");
props.put("io-impl", "org.apache.iceberg.aws.s3.S3FileIO");
RESTCatalog icebergCatalog = new RESTCatalog();
icebergCatalog.initialize(catalogName, props);配置項說明
配置項 | 說明 | 樣本 |
| DLF Iceberg REST 服務地址(VPC 內網),格式為 |
|
| DLF 資料目錄(Catalog)名稱。 |
|
| 認證類型,固定值: |
|
| 委託認證類型,固定值: |
|
| DLF 執行個體所在的 Region ID。 |
|
| 簽名服務名稱,固定值: |
|
| 訪問 DLF 的 AccessKey ID。 | |
| 訪問 DLF 的 AccessKey Secret。 | |
| 檔案 IO 實作類別,使用 Iceberg 社區標準 S3FileIO,固定值: |
|
開源 Spark 訪問
開源 Apache Spark 直接使用 Iceberg 社區 Runtime 即可讀寫 DLF Iceberg 表。Catalog 通過標準 Iceberg REST 協議通訊(AWS SigV4 簽名認證),資料面通過社區 S3FileIO 以 AWS S3 協議直連 OSS S3 相容端點。所有依賴均來自 Maven Central。
啟動 Spark SQL
spark-sql \
--packages org.apache.iceberg:iceberg-spark-runtime-3.5_2.12:1.11.0,org.apache.iceberg:iceberg-aws-bundle:1.11.0 \
--conf spark.sql.extensions=org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions \
--conf spark.sql.catalog.dlf=org.apache.iceberg.spark.SparkCatalog \
--conf spark.sql.catalog.dlf.type=rest \
--conf spark.sql.catalog.dlf.uri=http://cn-hangzhou-vpc.dlf.aliyuncs.com/iceberg \
--conf spark.sql.catalog.dlf.warehouse=iceberg_table_test \
--conf spark.sql.catalog.dlf.rest.auth.type=sigv4 \
--conf spark.sql.catalog.dlf.rest.auth.sigv4.delegate-auth-type=none \
--conf spark.sql.catalog.dlf.rest.signing-region=cn-hangzhou \
--conf spark.sql.catalog.dlf.rest.signing-name=DlfNext \
--conf spark.sql.catalog.dlf.rest.access-key-id=xxx \
--conf spark.sql.catalog.dlf.rest.secret-access-key=yyy \
--conf spark.sql.catalog.dlf.io-impl=org.apache.iceberg.aws.s3.S3FileIOSpark SQL 讀寫樣本
配置完成後,使用標準 Spark SQL 文法讀寫 Iceberg 表:
CREATE DATABASE IF NOT EXISTS dlf.demo_db;
CREATE TABLE dlf.demo_db.demo_tbl (id BIGINT, name STRING) USING iceberg;
INSERT INTO dlf.demo_db.demo_tbl VALUES (1, 'hello'), (2, 'world');
SELECT * FROM dlf.demo_db.demo_tbl;開源 Flink 訪問
開源 Apache Flink 使用 Iceberg 社區 Runtime 訪問 DLF Iceberg 表。將 iceberg-flink-runtime-1.20-1.11.0.jar 與 iceberg-aws-bundle-1.11.0.jar(均為 Maven Central 社區構件)放入 Flink 的 lib/ 目錄,然後在 SQL Client 中註冊 Catalog。
註冊 Catalog
CREATE CATALOG dlf WITH (
'type' = 'iceberg',
'catalog-type' = 'rest',
'uri' = 'http://cn-hangzhou-vpc.dlf.aliyuncs.com/iceberg',
'warehouse' = 'iceberg_table_test',
'rest.auth.type' = 'sigv4',
'rest.auth.sigv4.delegate-auth-type' = 'none',
'rest.signing-region' = 'cn-hangzhou',
'rest.signing-name' = 'DlfNext',
'rest.access-key-id' = 'xxx',
'rest.secret-access-key' = 'yyy',
'io-impl' = 'org.apache.iceberg.aws.s3.S3FileIO'
);Flink SQL 讀寫樣本
USE CATALOG dlf;
CREATE DATABASE IF NOT EXISTS demo_db;
CREATE TABLE demo_db.demo_tbl (id BIGINT, name STRING);
INSERT INTO demo_db.demo_tbl VALUES (1, 'hello'), (2, 'world');
SELECT * FROM demo_db.demo_tbl;