本文介绍如何通过 MaxCompute Managed Iceberg Table的REST Catalog API,在 Spark 环境中按照完全兼容开源的方式,读取Iceberg 表的元数据和数据。
前提条件
确保当前操作项目已经通过试用表单,可以正常使用MaxCompute管理的Iceberg表(beta)。
将以下 JAR 包上传到 Spark 环境:
iceberg-odps-1.12.0-SNAPSHOT.jar
内含 OdpsAuthManager 的实现,必须上传。
iceberg-spark-runtime-3.5_2.12-1.10.1.jar
内含 Iceberg 在 Spark 上运行所需的所有依赖,若已有内置jar,不需额外上传。
内含 Iceberg 与 S3/OSS 兼容服务交互所需的依赖,若已有内置jar,不需额外上传。
Spark 3 服务需升级到 JDK 17。若使用阿里云 E-MapReduce(EMR)实例,可参考Spark3使用JDK 11,并将其中的 java-11 替换为 java-17。
配置并启动 Spark SQL
在 Spark 客户端节点(例如 EMR 集群的 Master 节点或本地已配置 Spark 的环境)执行以下命令启动 spark-sql:
启动spark-sql
export AWS_REGION=cn-shanghai
spark-sql \
--jars /root/iceberg/iceberg-spark-runtime-3.5_2.12-1.10.1.jar,/root/iceberg/iceberg-aws-bundle-1.10.1.jar,/root/iceberg/iceberg-odps-1.12.0-SNAPSHOT.jar \
--conf spark.driver.extraClassPath=/root/iceberg/* \
--conf spark.executor.extraClassPath=/root/iceberg/* \
--conf spark.sql.extensions=org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions \
--conf spark.sql.catalog.odps=org.apache.iceberg.spark.SparkCatalog \
--conf spark.sql.catalog.odps.catalog-impl=org.apache.iceberg.rest.RESTCatalog \
--conf spark.sql.catalog.odps.uri=https://catalogapi.cn-shanghai.maxcompute.aliyun-inc.com/api/iceberg/v1alpha/restcatalog/ \
--conf spark.sql.catalog.odps.warehouse=projects/project_name \
--conf spark.sql.catalog.odps.io-impl=org.apache.iceberg.aws.s3.S3FileIO \
--conf spark.sql.catalog.odps.rest.auth.type=org.apache.iceberg.odps.auth.OdpsAuthManager \
--conf spark.sql.catalog.odps.odps.auth.access-key-id=*** \
--conf spark.sql.catalog.odps.odps.auth.access-key-secret=*** \
--conf spark.sql.catalog.odps.s3.endpoint=https://oss-cn-shanghai.aliyuncs.com \
--conf spark.sql.catalog.odps.client.region=cn-shanghai \
--conf spark.sql.catalogImplementation=in-memory \
--conf spark.hadoop.javax.jdo.option.ConnectionURL="jdbc:derby:/tmp/metastore_db;create=true"运行spark-sql
启动 spark-sql 后,按项目是否启用 Schema 层级模式,选择对应语法访问 MaxCompute Iceberg 表。
项目支持Schema层级模式
支持Schema层级模式是指模式打开后,项目内支持表-schema-项目 的层级关系,支持schema操作,schema层级映射为Iceberg catalog的namespace层级,Rest Catalog API的语法格式为catalog.namespace.table,详情可见Schema操作。
列出项目内的namespace,只有三层模型的项目支持此语法。
show namespaces in odps;列出namespace下的表。
-- 列出名为default的namespace下的表。 show tables in odps.default; -- 列出名为test01的namespace下的表,test01可替换为其他schema名称。 show tables in odps.test01;查询namespace下的表。
-- 查询名为default的namespace下的表。 select * from odps.default.table_name; -- 查询名为test01的namespace下的表,test01可替换为其他schema名称。 select * from odps.test01.table_name;查看表详情。
-- 查看名为default的namespace下的表详情。 desc odps.default.table_name; -- 查询名为test01的namespace下的表详情,test01可替换为其他schema名称。 desc odps.test01.table_name;
项目不支持Schema层级模式
不支持Schema层级模式是指当前项目内表-项目 的层级关系,不支持schema相关操作,Rest Catalog Api的语法格式为project.default.table,详情可见Schema操作。
列出namespace下的表。
show tables in odps.default;查询namespace下的表。
select * from odps.default.table_name;查看表详情。
desc odps.default.table_name;
使用示例
以下以 EMR、RamRole 用户及支持schema层级的项目为例,展示完整操作流程。
准备数据
创建connection
登录MaxCompute控制台,在左上角选择地域。
在左侧导航栏,选择。
在数据湖连接(CONNECTION)页面,单击创建数据湖连接。
在弹出的创建数据湖连接对话框,填写如下参数,然后单击确定完成创建数据湖连接。
参数名称
说明
数据湖连接名称
数据湖连接名称,在租户内命名唯一。
RAMRoleARN
选择RAM Role中具有访问OSS权限的RAMRoleARN信息。
可以创建和填写自定义角色的RAMRoleARN信息,创建详情请参见访问外部数据源授权方案。
数据湖连接描述
数据湖连接描述。
在数据湖连接(CONNECTION)页面,单击要授予其他用户使用的CONNECTION对应的操作列的新增授权。
在弹出的数据湖连接授权对话框,添加需要授权的用户,单击确定完成授权。
创建manage iceberg table,并写入测试数据
创建 MaxCompute 管理的 Iceberg 表并写入测试数据。
SET odps.namespace.schema=true; CREATE ICEBERG TABLE default.mc_iceberg_table_rest_catalog ( emp_id BIGINT COMMENT '员工ID', emp_name STRING COMMENT '员工姓名', department STRING COMMENT '部门', salary DOUBLE COMMENT '薪资', hire_date DATE COMMENT '入职日期', ds STRING COMMENT '分区日期' ) PARTITIONED BY (ds) WITH CONNECTION <YOUR CONNECTION> OPTIONS( location='oss://<oss bucket>/rest_catalog/' ); -- 写表 SET odps.namespace.schema=true; INSERT INTO TABLE default.mc_iceberg_table_rest_catalog VALUES (1001, '张三', '技术部', 15000.00, date '2020-03-15','2024-01-01'); INSERT INTO TABLE default.mc_iceberg_table_rest_catalog VALUES (1002, '李四', '技术部', 18000.00, date '2019-07-22','2024-01-01'); INSERT INTO TABLE default.mc_iceberg_table_rest_catalog VALUES (1003, '王五', '产品部', 16000.00, date '2021-01-10','2024-01-01'); -- 确认SQL是否可读表 SET odps.namespace.schema=true; SELECT * FROM default.mc_iceberg_table_rest_catalog; +------------+------------+------------+------------+------------+------------+ | emp_id | emp_name | department | salary | hire_date | ds | +------------+------------+------------+------------+------------+------------+ | 1003 | 王五 | 产品部 | 16000.0 | 2021-01-10 | 2024-01-01 | | 1001 | 张三 | 技术部 | 15000.0 | 2020-03-15 | 2024-01-01 | | 1002 | 李四 | 技术部 | 18000.0 | 2019-07-22 | 2024-01-01 | +------------+------------+------------+------------+------------+------------+
启动spark-sql
以下命令在已启用 VPC 内网访问的 EMR 集群中启动 spark-sql,使用RamRole用户访问 MaxCompute 已开通三层模型的项目。若是普通RamUser(子账号),需要去掉spark.sql.catalog.odps.odps.auth.sts-token参数。
export AWS_REGION=cn-shanghai
spark-sql \
--jars /root/iceberg/iceberg-spark-runtime-3.5_2.12-1.10.1.jar,/root/iceberg/iceberg-aws-bundle-1.10.1.jar,/root/iceberg/iceberg-odps-1.12.0-SNAPSHOT.jar \
--conf spark.driver.extraClassPath=/root/iceberg/* \
--conf spark.executor.extraClassPath=/root/iceberg/* \
--conf spark.sql.extensions=org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions \
--conf spark.sql.catalog.odps=org.apache.iceberg.spark.SparkCatalog \
--conf spark.sql.catalog.odps.catalog-impl=org.apache.iceberg.rest.RESTCatalog \
--conf spark.sql.catalog.odps.uri=https://catalogapi.cn-shanghai-vpc.maxcompute.aliyun-inc.com/api/iceberg/v1alpha/restcatalog/ \
--conf spark.sql.catalog.odps.warehouse=projects/<projectname> \
--conf spark.sql.catalog.odps.io-impl=org.apache.iceberg.aws.s3.S3FileIO \
--conf spark.sql.catalog.odps.rest.auth.type=org.apache.iceberg.odps.auth.OdpsAuthManager \
--conf spark.sql.catalog.odps.odps.auth.access-key-id=STS.*** \
--conf spark.sql.catalog.odps.odps.auth.access-key-secret=*** \
--conf spark.sql.catalog.odps.odps.auth.sts-token=*** \
--conf spark.sql.catalog.odps.s3.endpoint=https://oss-cn-shanghai-internal.aliyuncs.com \
--conf spark.sql.catalog.odps.client.region=cn-shanghai \
--conf spark.sql.catalogImplementation=in-memory \
--conf spark.hadoop.javax.jdo.option.ConnectionURL="jdbc:derby:/tmp/metastore_db;create=true"运行spark-sql
列出项目内的namespace
spark-sql (default)> show namespaces in odps; default test01 Time taken: 1.676 seconds, Fetched 2 row(s)列出namespace下的表
spark-sql (default)> show tables in odps.default; mc_iceberg_table_test1 mc_iceberg_table_test2 mc_iceberg_table_test3 mc_iceberg_table_rest_catalog Time taken: 0.276 seconds, Fetched 4 row(s)查询namespace下的表
spark-sql (default)> select * from odps.default.mc_iceberg_table_rest_catalog; 1003 王五 产品部 16000.0 2021-01-10 2024-01-01 1002 李四 技术部 18000.0 2019-07-22 2024-01-01 1001 张三 技术部 15000.0 2020-03-15 2024-01-01 Time taken: 5.403 seconds, Fetched 3 row(s)查看表详情
spark-sql (default)> desc odps.test01.mc_iceberg_table_rest_catalog; emp_id bigint 员工ID emp_name string 员工姓名 department string 部门 salary double 薪资 hire_date date 入职日期 ds string 分区日期 # Partition Information # col_name data_type comment ds string 分区日期 Time taken: 0.78 seconds, Fetched 9 row(s)