Apache Iceberg REST Catalog API

更新时间:
复制 MD 格式

本文介绍如何通过 MaxCompute Managed Iceberg TableREST Catalog API,在 Spark 环境中按照完全兼容开源的方式,读取Iceberg 表的元数据和数据。

前提条件

配置并启动 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"

单击查看详细参数说明

参数

说明

示例

环境变量参数

export AWS_REGION

区分部署环境:

  • 本地 Spark 必须设置;

  • EMR 集群可省略。

固定值。cn-shanghai

jar包参数

--jars

启动spark-sql依赖的jar包。

需要提前上传到spark集群中,多个jar包以英文逗号分隔。

/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

Spark环境为EMR实例,为避免iceberg jar包冲突,需要加此参数。值为iceberg jar包的路径。

iceberg jar包在/root/iceberg/路径下,则该值为/root/iceberg/*

spark.executor.extraClassPath

Spark环境为EMR实例,为避免iceberg jar包冲突,需要加此参数。值为iceberg jar包的路径。

iceberg jar包在/root/iceberg/路径下,则该值为/root/iceberg/*

spark.sql.extensions

SparkSession 的扩展注册参数。

固定值。org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions

spark.sql.catalog.odps

注册一个名为 odps 的自定义 catalog,之后就可以在 Spark SQL 里直接用 odps 作为 catalog 名访问表。可以把 odps 换成自定义名字,同时也要把下面参数里的 odps 全部换掉。

固定值。org.apache.iceberg.spark.SparkCatalog

spark.sql.catalog.odps.catalog-impl

Catalog 的实现。

固定值。org.apache.iceberg.rest.RESTCatalog

spark.sql.catalog.odps.uri

MaxCompute Rest Catalog endpoint

  • Spark集群在本地或EMR集群开了公网,使用公网地址:https://catalogapi.cn-shanghai.maxcompute.aliyun-inc.com/api/iceberg/v1alpha/restcatalog/

  • Spark集群在EMR未开公网,使用VPC地址:https://catalogapi.cn-shanghai-vpc.maxcompute.aliyun-inc.com/api/iceberg/v1alpha/restcatalog/

spark.sql.catalog.odps.warehouse

Iceberg表所在的MaxCompute 项目名称。

spark.sql.catalog.odps.io-impl

数据 IO 实现。

固定值。org.apache.iceberg.aws.s3.S3FileIO

spark.sql.catalog.odps.rest.auth.type

身份认证的实现。

固定值。org.apache.iceberg.odps.auth.OdpsAuthManager。这是 MaxCompute 的身份认证实现。

spark.sql.catalog.odps.odps.auth.access-key-id

账号AccessId。

spark.sql.catalog.odps.odps.auth.access-key-secret

账号AccessKey。

spark.sql.catalog.odps.odps.auth.sts-token

账号sts-token。只有当使用RamRole用户操作时,需要填写此参数,主账号或RamUser操作时,不需要填写。

spark.sql.catalog.odps.s3.endpoint

数据链路 IO 使用的 OSS endpoint

  • Spark集群在本地或EMR集群开了公网,使用公网地址:https://oss-cn-shanghai.aliyuncs.com

  • Spark集群在EMR未开公网,使用vpc地址:https://oss-cn-shanghai-internal.aliyuncs.com

  • 内网可达的环境下建议优先使用内网域名,公网访问可能产生OSS侧额外的费用。

spark.sql.catalog.odps.client.region

Spark环境为EMR实例,需要加此配置项。

固定值。cn-shanghai

spark.sql.catalogImplementation

Spark环境为EMR实例,需要加此配置项。

固定值。in-memory

spark.hadoop.javax.jdo.option.ConnectionURL

Spark环境为EMR实例,需要加此配置项。

固定值。jdbc:derby:/tmp/metastore_db;create=true

运行spark-sql

启动 spark-sql 后,按项目是否启用 Schema 层级模式,选择对应语法访问 MaxCompute Iceberg 表。

项目支持Schema层级模式

支持Schema层级模式是指模式打开后,项目内支持表-schema-项目 的层级关系,支持schema操作,schema层级映射为Iceberg catalognamespace层级,Rest Catalog API的语法格式为catalog.namespace.table,详情可见Schema操作

  1. 列出项目内的namespace,只有三层模型的项目支持此语法。

    show namespaces in odps;
  2. 列出namespace下的表。

    -- 列出名为defaultnamespace下的表。
    show tables in odps.default;
    
    -- 列出名为test01namespace下的表,test01可替换为其他schema名称。
    show tables in odps.test01;
  3. 查询namespace下的表。

    -- 查询名为defaultnamespace下的表。
    select * from odps.default.table_name;
    
    -- 查询名为test01namespace下的表,test01可替换为其他schema名称。
    select * from odps.test01.table_name;
  4. 查看表详情。

    -- 查看名为defaultnamespace下的表详情。
    desc odps.default.table_name;
    
    -- 查询名为test01namespace下的表详情,test01可替换为其他schema名称。
    desc odps.test01.table_name;

项目不支持Schema层级模式

不支持Schema层级模式是指当前项目内表-项目 的层级关系,不支持schema相关操作,Rest Catalog Api的语法格式为project.default.table,详情可见Schema操作

  1. 列出namespace下的表。

    show tables in odps.default;
  2. 查询namespace下的表。

    select * from odps.default.table_name;
  3. 查看表详情。

    desc odps.default.table_name;

使用示例

以下以 EMR、RamRole 用户及支持schema层级的项目为例,展示完整操作流程。

准备数据

  1. 创建connection

    1. 登录MaxCompute控制台,在左上角选择地域。

    2. 在左侧导航栏,选择MaxLake > 数据湖连接

    3. 数据湖连接(CONNECTION)页面,单击创建数据湖连接

    4. 在弹出的创建数据湖连接对话框,填写如下参数,然后单击确定完成创建数据湖连接。

      参数名称

      说明

      数据湖连接名称

      数据湖连接名称,在租户内命名唯一。

      RAMRoleARN

      选择RAM Role中具有访问OSS权限的RAMRoleARN信息。

      可以创建和填写自定义角色的RAMRoleARN信息,创建详情请参见访问外部数据源授权方案

      数据湖连接描述

      数据湖连接描述。

    5. 数据湖连接(CONNECTION)页面,单击要授予其他用户使用的CONNECTION对应的操作列的新增授权

      在弹出的数据湖连接授权对话框,添加需要授权的用户,单击确定完成授权。

  2. 创建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

  1. 列出项目内的namespace

    spark-sql (default)> show namespaces in odps;
    default
    test01
    Time taken: 1.676 seconds, Fetched 2 row(s)
  2. 列出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)
  3. 查询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)
  4. 查看表详情

    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)