本文为您介绍在 EMR on ECS Spark 环境中,如何使用开源 Iceberg Spark Runtime 通过 Iceberg REST 访问 DLF Iceberg 表。
前提条件
版本要求:已创建 EMR-5.12.0 及以上版本的集群,选择 Spark3 组件。Spark 使用 JDK 17(Iceberg 1.11.0 要求,参照 Spark3使用JDK 11 的同样方式配置)。
地域要求:EMR 集群与 DLF 处于同一地域,且集群所在 VPC 已加入 DLF 白名单。
权限要求:拥有访问 DLF 的 AccessKey,对应 RAM 用户已被授予目标 Catalog 的数据权限。详情请参见数据授权管理。
依赖引入
仅需以下两个 Apache Iceberg 社区构件(Maven Central):
iceberg-spark-runtime-3.5_2.12(1.11.0 及以上版本)iceberg-aws-bundle(1.11.0 及以上版本)
可通过 --packages 在提交时自动拉取(见下方示例);离线集群也可提前下载后放入 $SPARK_HOME/jars。
使用示例
配置 Catalog 连接
在 Terminal 中执行 spark-sql 命令,注意替换对应参数。
spark-sql \
--master local \
--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.iceberg_catalog=org.apache.iceberg.spark.SparkCatalog \
--conf spark.sql.catalog.iceberg_catalog.catalog-impl=org.apache.iceberg.rest.RESTCatalog \
--conf spark.sql.catalog.iceberg_catalog.uri=http://${regionID}-vpc.dlf.aliyuncs.com/iceberg \
--conf spark.sql.catalog.iceberg_catalog.warehouse=${catalogName} \
--conf spark.sql.catalog.iceberg_catalog.io-impl=org.apache.iceberg.aws.s3.S3FileIO \
--conf spark.sql.catalog.iceberg_catalog.rest.auth.type=sigv4 \
--conf spark.sql.catalog.iceberg_catalog.rest.auth.sigv4.delegate-auth-type=none \
--conf spark.sql.catalog.iceberg_catalog.rest.signing-region=${regionID} \
--conf spark.sql.catalog.iceberg_catalog.rest.signing-name=DlfNext \
--conf spark.sql.catalog.iceberg_catalog.rest.access-key-id=${AccessKeyId} \
--conf spark.sql.catalog.iceberg_catalog.rest.secret-access-key=${AccessKeySecret}配置项说明如下表所示:
配置项 | 说明 | 示例 |
| DLF Iceberg REST 服务地址(VPC 内网),格式为 |
|
| 数据目录名 |
|
| 设置为固定值: |
|
| Iceberg 社区标准实现,设置为固定值: |
|
| auth 的类型,设置为固定值: |
|
| 设置为固定值: |
|
| DLF 的 Region ID |
|
| 设置为固定值: |
|
| 访问 DLF 的 AccessKeyId | |
| 访问 DLF 的 AccessKeySecret |
读写 DLF Iceberg
启动后即可用标准 Spark SQL 读写:
CREATE DATABASE IF NOT EXISTS iceberg_catalog.db;
CREATE TABLE iceberg_catalog.db.iceberg_tbl (id BIGINT, name STRING) USING iceberg;
INSERT INTO iceberg_catalog.db.iceberg_tbl VALUES (1, 'hello'), (2, 'world');
SELECT * FROM iceberg_catalog.db.iceberg_tbl;