このトピックでは、オープンソースの Apache Iceberg Java API、Spark、および Flink を使用して、DLF Iceberg REST カタログにアクセスする方法について説明します。AWS SigV4 署名を使用した標準の Iceberg REST プロトコルにより、テーブルメタデータの管理や Iceberg テーブルの読み書きが可能です。データプレーンは、コミュニティ標準の S3FileIO を使用して OSS S3 互換エンドポイントに直接接続します。
前提条件
ランタイム環境では JDK 17 以降が必要です (Iceberg 1.11.0 で必須)。
次の Apache Iceberg コミュニティの依存関係 (バージョン 1.11.0) を Maven プロジェクトに追加してください。すべての依存関係は 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 を使用した接続
コミュニティ標準の RESTCatalog クラスを使用して、DLF Iceberg REST カタログに接続します。初期化時に次の構成パラメータを渡します。
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 カタログ名。 |
|
| 認証タイプ。 |
|
| 委任認証タイプ。 |
|
| DLF インスタンスのリージョン ID。 |
|
| 署名サービス名。 |
|
| DLF へのアクセスに使用する AccessKey ID。 | |
| DLF へのアクセスに使用する AccessKey Secret。 | |
| File IO 実装クラス。コミュニティ標準の S3FileIO を使用します。 |
|
Apache Spark を使用したアクセス
オープンソースの Apache Spark は、Iceberg コミュニティランタイムを使用して DLF Iceberg テーブルを直接読み書きできます。カタログは AWS SigV4 署名を使用した標準の Iceberg REST プロトコルで通信し、データプレーンはコミュニティの S3FileIO を使用して、S3 互換 API 経由で 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;Apache Flink を使用したアクセス
オープンソースの Apache Flink は、Iceberg コミュニティランタイムを使用して DLF Iceberg テーブルにアクセスします。iceberg-flink-runtime-1.20-1.11.0.jar と iceberg-aws-bundle-1.11.0.jar (いずれも Maven Central のコミュニティアーティファクト) を Flink インストール先の lib/ ディレクトリに配置し、SQL クライアントでカタログを登録してください。
カタログの登録
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;