通过Spark SQL读写Iceberg外表

更新时间:
复制 MD 格式

云原生数据仓库 AnalyticDB MySQL 版支持通过Spark SQL创建、读写和管理Iceberg外表。Spark 3.5及以上版本默认内置Iceberg读写能力,无需额外配置即可使用。本文介绍如何选择Catalog模式管理Iceberg表元数据,以及如何完成Iceberg外表的创建、读写、查询和删除操作。

前提条件

Catalog模式说明

Spark读写Iceberg外表时,需要通过Catalog管理表的元数据。AnalyticDB for MySQL支持以下三种Catalog模式:

Catalog模式

说明

是否需要额外配置

适用场景

MDS Catalog(默认)

基于AnalyticDB for MySQL内置元数据服务(MDS)管理Iceberg表元数据。

不需要,开箱即用。

仅在AnalyticDB for MySQL内部使用Spark读写Iceberg表。

Hadoop Catalog

基于Hadoop文件系统管理Iceberg表元数据,元数据以文件形式存储在指定的OSS路径下。

需要配置spark.iceberg.warehouse

需要跨产品协同读写同一张表,例如EMR Spark、Flink等。

湖存储

基于AnalyticDB for MySQL湖存储管理Iceberg表,底层使用Hadoop Catalog,由系统自动管理存储路径。

需要开启spark.adb.lakehouse.enabled

湖仓版集群使用湖存储功能。

如何选择Catalog模式

  • 仅在AnalyticDB for MySQL内部使用:三种模式均可,推荐使用湖存储模式(若已开通湖存储)或默认的MDS Catalog模式,配置最简单。

  • 需要跨产品协同读写同一张表(例如EMR Spark写入后由AnalyticDB for MySQL查询,或Flink写入后由Spark查询):建议使用Hadoop Catalog或湖存储模式。不同Catalog模式的元数据相互独立,跨模式混合写入同一张表会导致元数据不一致,详情请参见注意事项

表名引用方式

AnalyticDB for MySQL支持以下两种方式引用Iceberg表:

  • 两层结构(推荐)数据库名.表名,例如my_db.my_table。系统会自动识别表类型并路由到对应的Iceberg Catalog,无需指定Catalog名称。

  • 三层结构Catalog名.数据库名.表名,例如iceberg.my_db.my_table。适用于需要明确指定Catalog的场景。

Spark 3.5及以上版本中,系统自动注册了以下Iceberg Catalog名称:

Catalog名称

说明

adb_lakehouse_prod

默认的Iceberg Catalog名称。

iceberg

Iceberg Catalog的简便别名,与adb_lakehouse_prod等价。

示例(以下三种写法等价

SELECT * FROM my_db.my_table;
SELECT * FROM iceberg.my_db.my_table;
SELECT * FROM adb_lakehouse_prod.my_db.my_table;

操作步骤

步骤一:进入数据开发

  1. 登录云原生数据仓库AnalyticDB MySQL控制台,在左上角选择集群所在地域。在左侧导航栏,单击集群列表,然后单击目标集群ID。

  2. 配置Spark参数。

    根据所选的Catalog模式,需要配置不同模式的SET参数:

    Catalog模式

    SET参数

    MDS Catalog(默认)

    SET spark.adb.version=3.5;

    Hadoop Catalog

    SET spark.adb.version=3.5;
    SET spark.iceberg.warehouse=oss://testBucketName/iceberg/;

    湖存储

    SET spark.adb.lakehouse.enabled=true;

    SET参数的配置方式因资源组类型而异:

    • Job型资源组 :直接将SET语句写在SQL语句前面即可。如本文后续SQL示例所示。

    • Interactive型资源组 :需要将Spark参数预先配置到资源组中。配置完成后,SQL语句中无需再添加对应的SET语句。

      1. 在左侧导航栏单击集群管理 > 资源管理

      2. 在目标资源组所在行操作列,单击修改

      3. 修改资源组面板中更新Spark 配置,然后单击确定

        当资源组状态变为运行中时,修改生效。

  3. 在左侧导航栏,单击作业开发 > SQL开发

  4. SQLConsole窗口,选择Spark引擎和资源组(Job型资源组或Spark引擎的Interactive型资源组)。

步骤二:创建外库与Iceberg外表

请根据所选的Catalog模式,按照对应方式创建外库和Iceberg外表。

MDS Catalog模式(默认)

Spark 3.5及以上版本默认内置Iceberg Catalog支持。MDS Catalog模式下,元数据由AnalyticDB for MySQL MDS自动管理,无需配置spark.iceberg.warehouseCatalog参数,但需要指定数据的存储位置(LOCATION)。

  1. 创建外库

    • 您可以在创建数据库时通过LOCATION指定默认的数据存储路径。指定后,该数据库下所有表将自动继承此路径。

      SET spark.adb.version=3.5;
      
      CREATE DATABASE adb_iceberg_db LOCATION 'oss://your-bucket/iceberg/';
    • 您也可以创建不带LOCATION的数据库。此时需要在建表时为每张表单独指定LOCATION

      CREATE DATABASE adb_iceberg_db;
  2. 创建Iceberg外表

    SET spark.adb.version=3.5;
    
    -- 若数据库已指定LOCATION,建表时无需再指定
    CREATE TABLE adb_iceberg_db.test_iceberg_tbl (
      `id` int,
      `name` string,
      `age` int
    ) USING iceberg
    PARTITIONED BY (age);
    
    -- 若数据库未指定LOCATION,建表时需通过LOCATION指定数据存储路径
    CREATE TABLE adb_iceberg_db.test_iceberg_tbl (
      `id` int,
      `name` string,
      `age` int
    ) USING iceberg
    PARTITIONED BY (age)
    LOCATION 'oss://your-bucket/iceberg/test_iceberg_tbl';

Hadoop Catalog模式

如需与其他引擎或产品(如EMR Spark、Flink等)共同读写同一张Iceberg表,建议使用Hadoop Catalog模式并指定OSS warehouse路径,具体配置请参见SET参数配置

  1. 创建外库前的检查

    如需使用已有数据库,请通过DESCRIBE DATABASE查看数据库DDL,确认满足以下条件之一:

    • DDL中未指定Location

    • DDL中指定了Location,且Catalog参数值为mix

    如不满足上述条件,请新建数据库。

  2. 创建外库

    CREATE DATABASE adb_external_db_iceberg;
    重要

    如果数据库指定了Location,请确保spark.iceberg.warehouse参数指定的OSS路径前缀与Location保持一致。

  3. 创建Iceberg外表

    SET spark.adb.version=3.5;                                  -- 指定Spark引擎大版本
    SET spark.iceberg.warehouse=oss://testBucketName/iceberg/;   -- Iceberg外表元数据与数据文件的存储路径
    
    CREATE TABLE adb_external_db_iceberg.test_iceberg_tbl (
      `id` int,
      `name` string,
      `age` int
    ) USING iceberg
    PARTITIONED BY (age);

湖存储模式

湖仓版集群支持通过湖存储模式创建和管理Iceberg表。湖存储模式底层使用Hadoop Catalog,由AnalyticDB for MySQL自动管理数据存储路径。

  1. 创建外库

    CREATE DATABASE adb_external_db_iceberg
      WITH DBPROPERTIES (
        'adb_lake_bucket' = 'adb-lake-cn-shanghai-6gml****'
      );

    adb_lake_bucket:非必须。指定湖存储表数据的存储位置。在数据库级别指定后,该数据库下所有表默认使用此湖存储空间。您也可以在建表时单独指定,表级别配置优先于数据库级别。

  2. 创建Iceberg外表

    SET spark.adb.lakehouse.enabled=true;                       -- 开启湖存储
    
    CREATE TABLE adb_external_db_iceberg.test_iceberg_tbl (
      `id` int,
      `name` string,
      `age` int
    ) USING iceberg
    PARTITIONED BY (age)
    TBLPROPERTIES (
      'adb_lake_bucket' = 'adb-lake-cn-shanghai-6gml****'
    );

步骤三:写入或删除Iceberg外表数据

以下示例以MDS Catalog模式为主进行说明。

写入数据

  • INSERT INTO(追加写入)

    SET spark.adb.version=3.5;
    
    INSERT INTO adb_iceberg_db.test_iceberg_tbl
    VALUES (1, 'lisa', 10), (2, 'jams', 20);
  • INSERT OVERWRITE(全量覆盖写入)

    INSERT OVERWRITE adb_iceberg_db.test_iceberg_tbl
    VALUES (1, 'lisa', 10), (2, 'jams', 30);
  • INSERT OVERWRITE静态分区写入

    仅覆盖指定分区的数据,其他分区不受影响。

    INSERT OVERWRITE adb_iceberg_db.test_iceberg_tbl
      PARTITION(age=10)
    VALUES (1, 'anna');
  • INSERT OVERWRITE动态分区写入

    仅覆盖写入数据所涉及的分区,其他分区不受影响。

    重要

    如未设置spark.sql.sources.partitionOverwriteMode=dynamic,INSERT OVERWRITE将覆盖全表数据,而非仅覆盖目标分区。

    SET spark.sql.sources.partitionOverwriteMode=dynamic;        -- 开启动态分区覆盖模式
    
    INSERT OVERWRITE adb_iceberg_db.test_iceberg_tbl
      PARTITION(age)
    VALUES (1, 'bom', 10);
  • UPDATE(条件更新)

    UPDATE adb_iceberg_db.test_iceberg_tbl
    SET name = 'box'
    WHERE id = 2;

删除数据

DELETE FROM adb_iceberg_db.test_iceberg_tbl
WHERE id = 1;

DELETE FROM adb_iceberg_db.test_iceberg_tbl
WHERE age = 20;

步骤四:查询Iceberg外表数据

SET spark.adb.version=3.5;

SELECT * FROM adb_iceberg_db.test_iceberg_tbl;

返回示例

+---+----+---+
|id |name|age|
+---+----+---+
|2  |box|30 |
+---+----+---+

步骤五:删除Iceberg外表

MDS Catalog模式

SET spark.adb.version=3.5;

-- 仅删除表的元数据,不删除底层数据文件
DROP TABLE adb_iceberg_db.test_iceberg_tbl;

-- 同时删除表的元数据和底层数据文件
DROP TABLE adb_iceberg_db.test_iceberg_tbl PURGE;

Hadoop Catalog模式

SET spark.adb.version=3.5;

DROP TABLE adb_iceberg_db.test_iceberg_tbl;

湖存储模式

SET spark.adb.lakehouse.enabled=true;

DROP TABLE adb_external_db_iceberg.test_iceberg_tbl;
说明

湖存储模式下,DROP TABLE仅删除表的元数据,不支持通过PURGE关键字删除底层数据文件,数据会由AnalyticDB for MySQL表回收站延迟删除。

高级功能:元数据查询与存储过程

Iceberg提供了元数据表查询和存储过程(Procedure)能力,用于表维护和元数据分析。

说明

元数据表查询和存储过程的调用与普通读写遵循相同的Catalog配置规则。使用两层结构(db.table)时,系统会自动路由到对应Catalog,无需额外配置;使用三层结构时需指定Catalog名称。对于通过Hadoop Catalog或湖存储模式创建的表,还需配置对应的SET参数(如spark.iceberg.warehousespark.adb.lakehouse.enabled)。

查询元数据表

Iceberg为每张表自动维护了多种元数据表,您可以通过表名.元数据表名的方式进行查询。

SET spark.adb.version=3.5;

-- 查看表的快照信息(快照ID、时间戳、操作类型等)
SELECT * FROM adb_iceberg_db.test_iceberg_tbl.snapshots;

-- 查看表的变更历史
SELECT * FROM adb_iceberg_db.test_iceberg_tbl.history;

-- 查看当前快照中的数据文件信息
SELECT * FROM adb_iceberg_db.test_iceberg_tbl.files;

-- 查看manifest条目(每个文件的详细引用信息)
SELECT * FROM adb_iceberg_db.test_iceberg_tbl.entries;

-- 查看分区统计信息(记录数、文件数等)
SELECT * FROM adb_iceberg_db.test_iceberg_tbl.partitions;

常用的元数据表包括snapshotshistoryfilesmanifestsentriespartitionsrefs等。此外,以all_为前缀的元数据表(如all_data_filesall_manifestsall_entries)可查看包含历史快照在内的全量信息。完整列表请参见Apache Iceberg官方文档

使用存储过程(Procedure)

您可以通过CALL语句调用存储过程进行表维护。

-- 两层结构:直接使用system命名空间
CALL system.procedure_name('数据库名.表名', 其他参数...);

-- 三层结构:需加上Catalog名前缀
CALL iceberg.system.procedure_name('数据库名.表名', 其他参数...);

以下是常用的表维护操作示例。

  • 合并小文件(rewrite_data_files)

    SET spark.adb.version=3.5;
    
    -- 使用默认binpack策略合并小文件
    CALL system.rewrite_data_files('adb_iceberg_db.test_iceberg_tbl');
    
    -- 使用sort策略并指定排序字段,仅对指定分区生效
    CALL system.rewrite_data_files(
      table => 'adb_iceberg_db.test_iceberg_tbl',
      strategy => 'sort',
      sort_order => 'id ASC NULLS FIRST',
      where => 'age = 10'
    );
  • 优化manifest文件(rewrite_manifests)

    CALL system.rewrite_manifests('adb_iceberg_db.test_iceberg_tbl');
  • 过期快照清理(expire_snapshots)

    CALL system.expire_snapshots(
      table => 'adb_iceberg_db.test_iceberg_tbl',
      older_than => TIMESTAMP '2024-01-01 00:00:00',
      retain_last => 5
    );
  • 快照回滚(rollback_to_snapshot)

    -- 先查询快照列表获取目标snapshot_id
    SELECT * FROM adb_iceberg_db.test_iceberg_tbl.snapshots;
    
    -- 回滚到指定快照
    CALL system.rollback_to_snapshot('adb_iceberg_db.test_iceberg_tbl', 1234567890);

除上述存储过程外,还支持remove_orphan_files(清理孤立文件)和rollback_to_timestamp(按时间回滚)等操作。完整列表请参见Apache Iceberg官方文档

AnalyticDB for MySQL增强存储过程

AnalyticDB for MySQL在标准Iceberg存储过程基础上,额外提供了数据归档和生命周期管理能力。

  • 数据归档(archive_files)

    将冷数据归档到低成本存储类型,降低存储费用。

    -- 归档30天前的数据到Archive存储(dry_run=true表示仅预览影响范围,不实际执行)
    CALL system.archive_files(
      table => 'adb_iceberg_db.test_iceberg_tbl',
      storage_type => 'Archive',
      start_days => 30,
      dry_run => true
    );
  • 检查归档状态(check_archive_status)

    查看各存储类型下的文件分布情况。

    CALL system.check_archive_status(
      table => 'adb_iceberg_db.test_iceberg_tbl',
      group_by_partition => true
    );
  • 数据解冻(restore_files)

    将已归档的数据恢复到标准存储类型。

    CALL system.restore_files(
      table => 'adb_iceberg_db.test_iceberg_tbl',
      where => 'age = 10'
    );
  • 分区生命周期管理(ttl)

    基于分区的最后更新时间,识别或删除过期分区。

    -- 查看超过90天未更新的分区(dry_run=true表示仅预览,不执行删除)
    CALL system.ttl(
      table => 'adb_iceberg_db.test_iceberg_tbl',
      expire_after_days => 90,
      dry_run => true
    );

注意事项

不同Catalog模式的元数据一致性

三种Catalog模式各自独立管理Iceberg表的元数据。在AnalyticDB for MySQL内部,MDS Catalog能够感知通过Hadoop Catalog或湖存储模式创建的表,因此在不添加额外配置的情况下也可以读取这些表。但这不意味着可以跨模式混合写入。

Catalog模式之间的跨模式读写行为

表的创建模式

通过MDS Catalog(默认)

通过Hadoop Catalog

通过湖存储

MDS Catalog创建的表

读写正常

无法识别

无法识别

Hadoop Catalog创建的表

可读取,不建议写入

读写正常

不建议混用

湖存储创建的表

可读取,不建议写入

可读写,但无法自动管理路径

读写正常

重要

请勿跨Catalog模式混合写入同一张表。虽然MDS Catalog可以读取Hadoop Catalog和湖存储模式创建的表,但如果通过MDS Catalog写入这些表,更新的元数据仅会记录在MDS中,不会同步到Hadoop Catalog的文件系统元数据。这将导致其他引擎读取到的数据与实际数据不一致。

建议对同一张表,始终使用创建该表时所用的Catalog模式进行读写操作。如需跨产品协同,请在建表时选择Hadoop Catalog或湖存储模式。

运行环境配置

  • 使用Iceberg相关功能时,Spark版本需为3.5及以上。如果集群默认Spark版本不是3.5,请在SQL语句开头添加SET spark.adb.version=3.5;

  • 根据所选的Catalog模式,需要配置不同模式的SET参数,不同资源组的配置方式,请参见SET参数配置

批处理与交互式模式

Spark SQL支持批处理和交互式两种运行模式。请确认所使用的资源组类型与运行模式相匹配。