API

API

师成师成· 更新于 2026-09-28· 阅读 14 分钟· 0 次阅读

登录后可跨设备保存划线和私人笔记登录

Iceberg Java API

表

Iceberg API 的主要用途是管理表元数据,例如模式(schema)、分区规范、元数据以及存储表数据的数据文件。

表元数据和相关操作通过 Table 接口访问,该接口用于返回表信息。

表元数据

Table 接口提供对表元数据的访问:

  • schema 返回当前表的模式
  • spec 返回当前表的分区规范
  • properties 返回键值对形式的属性映射
  • currentSnapshot 返回当前表的快照
  • snapshots 返回该表的所有有效快照
  • snapshot(id) 根据 ID 返回指定快照
  • location 返回表的根目录位置

表还提供 refresh 方法,用于将表更新到最新版本,并暴露以下辅助方法:

  • io 返回用于读写表文件的 FileIO
  • locationProvider 返回用于生成数据文件和元数据文件路径的 LocationProvider

扫描

文件级别

Iceberg 表扫描首先通过 newScan 创建 TableScan 对象。

TableScan scan = table.newScan();

要配置扫描,请在 TableScan 上调用 filter 和 select,以获取应用了这些更改的新 TableScan。

TableScan filteredScan = scan.filter(Expressions.equal("id", 5))

对配置方法的调用会创建一个新的 TableScan,从而保证每个 TableScan 都是不可变的,即使在多个线程间共享也不会发生意外变化。

配置好扫描后,planFiles、planTasks 和 schema 用于返回文件、任务以及读取投影。

TableScan scan = table.newScan()
    .filter(Expressions.equal("id", 5))
    .select("id", "data");

Schema projection = scan.schema();
Iterable<CombinedScanTask> tasks = scan.planTasks();

使用 asOfTime 或 useSnapshot 为时间旅行查询配置表快照。

行级

Iceberg 表扫描首先通过 IcebergGenerics.read 创建 ScanBuilder 对象。

ScanBuilder scanBuilder = IcebergGenerics.read(table)

要配置扫描,可在 ScanBuilder 上调用 where 和 select,以获取包含这些更改的新 ScanBuilder。

scanBuilder.where(Expressions.equal("id", 5))

配置好扫描后,调用 build 方法来执行扫描。build 返回 CloseableIterable<Record>。

CloseableIterable<Record> result = IcebergGenerics.read(table)
        .where(Expressions.lessThan("id", 5))
        .build();

其中 Record 是 iceberg-data 模块的 Iceberg 记录类型 org.apache.iceberg.data.Record。

更新操作

Table 还提供了用于更新表的操作。这些操作采用构建器模式,即 PendingUpdate,当调用 PendingUpdate#commit 时才会提交。

例如,更新表结构(schema)的过程是:调用 updateSchema,向构建器中添加更新内容,最后调用 commit 将待提交的更改提交到表中:

table.updateSchema()
    .addColumn("count", Types.LongType.get())
    .commit();

可用于更新表的操作有:

  • updateSchema -- 更新表结构
  • updateSpec -- 修改表的分区规范
  • updateStatistics -- 更新表的统计信息文件
  • updatePartitionStatistics -- 更新表中某个特定分区的统计信息
  • updateProperties -- 更新表属性
  • updateLocation -- 更新表的基础路径
  • expireSnapshots -- 用于从表中移除旧的快照
  • manageSnapshots -- 用于管理表的快照
  • newAppend -- 用于追加数据文件
  • newFastAppend -- 用于追加数据文件,不会压缩元数据
  • newOverwrite -- 用于追加数据文件并移除被覆盖的文件
  • newDelete -- 用于删除数据文件
  • newRewrite -- 用于重写数据文件;会用新版本替换现有文件
  • newRowDelta -- 用于删除或替换现有数据文件中的行
  • newTransaction -- 创建一个新的表级事务
  • rewriteManifests -- 通过对文件进行聚类来重写 manifest 数据,以加快扫描规划速度
  • replaceSortOrder -- 用新创建的排序顺序替换表的排序顺序
  • newReplacePartitions -- 用于用新数据动态覆盖表中的分区

事务

事务用于在单个原子操作中提交多项表变更。事务通过工厂方法(如 newAppend)创建各个操作,其使用方式与使用 Table 相同。由事务创建的操作会在调用 commitTransaction 时作为一组一起提交。

例如,在同一个事务中删除并追加一个文件:

Transaction t = table.newTransaction();

// commit operations to the transaction
t.newDelete().deleteFromRowFilter(filter).commit();
t.newAppend().appendFile(data).commit();

// commit all the changes to the table
t.commitTransaction();

类型

Iceberg 的数据类型位于 org.apache.iceberg.types 包中。

基本类型

基本类型的实例可以通过各类型类中的静态方法获取。没有参数的类型使用 get,而像 decimal 这样的类型则使用工厂方法:

Types.IntegerType.get()    // int
Types.DoubleType.get()     // double
Types.DecimalType.of(9, 2) // decimal(9, 2)

嵌套类型

结构体(struct)、映射(map)和列表(list)通过类型类中的工厂方法创建。

与结构体字段一样,映射的键或值以及列表的元素也会作为嵌套字段进行跟踪。嵌套字段会跟踪字段 ID 和可空性。

结构体字段使用 NestedField.optional 或 NestedField.required 创建。映射值和列表元素的可空性在映射和列表的工厂方法中设置。

// struct<1 id: int, 2 data: optional string>
StructType struct = Struct.of(
    Types.NestedField.required(1, "id", Types.IntegerType.get()),
    Types.NestedField.optional(2, "data", Types.StringType.get())
  )
// map<1 key: int, 2 value: optional string>
MapType map = MapType.ofOptional(
    1, 2,
    Types.IntegerType.get(),
    Types.StringType.get()
  )
// array<1 element: int>
ListType list = ListType.ofRequired(1, IntegerType.get());

表达式

Iceberg 的表达式用于配置表扫描。要创建表达式,请使用 Expressions 中的工厂方法。

支持的谓词表达式有:

  • isNull
  • notNull
  • equal
  • notEqual
  • lessThan
  • lessThanOrEqual
  • greaterThan
  • greaterThanOrEqual
  • in
  • notIn
  • startsWith
  • notStartsWith

支持的表达式运算有:

  • and
  • or
  • not

常量表达式有:

  • alwaysTrue
  • alwaysFalse

表达式绑定

表达式在创建时处于未绑定状态。在使用表达式之前,需要将其绑定到某个数据类型,以查找表达式名称所代表的字段 ID,并转换谓词字面量。

例如,在使用表达式 lessThan("x", 10) 之前,Iceberg 需要确定 "x" 指的是哪一列,并将 10 转换为该列的数据类型。

如果该表达式被绑定到类型 struct<1 x: long, 2 y: long> 或 struct<11 x: int, 12 y: int>。

表达式示例

table.newScan()
    .filter(Expressions.greaterThanOrEqual("x", 5))
    .filter(Expressions.lessThan("x", 10))

模块

Iceberg 表的支持按库模块组织:

  • iceberg-common 包含其他模块中使用的工具类
  • iceberg-api 包含公开的 Iceberg API,包括表达式、类型、表和操作
  • iceberg-arrow 是 Iceberg 类型系统的实现,用于以 Apache Arrow 作为内存数据格式来读取和写入存储在 Iceberg 表中的数据
  • iceberg-aws 包含 Iceberg API 的实现,用于访问存储在 AWS S3 上的表,以及/或者用于使用 AWS Glue 数据目录定义的表
  • iceberg-core 包含 Iceberg API 的实现以及对 Avro 数据文件的支持,处理引擎应当依赖此模块
  • iceberg-parquet 是一个可选模块,用于处理由 Parquet 文件支撑的表
  • iceberg-orc 是一个可选模块,用于处理由 ORC 文件支撑的表(实验性功能)
  • iceberg-hive-metastore 是由 Hive metastore Thrift 客户端支撑的 Iceberg 表实现

本项目 Iceberg 还包含用于为处理引擎及相关工具添加 Iceberg 支持的模块:

  • iceberg-spark 是 Iceberg 的 Spark Datasource V2 API 实现,针对每个 Spark 版本都有子模块(如需使用经过 shade 处理的版本,请使用 runtime jar)
  • iceberg-flink 是 Iceberg 的 Flink Table 和 DataStream API 实现(如需使用经过 shade 处理的版本,请使用 iceberg-flink-runtime)
  • iceberg-mr 是 Iceberg 的 MapReduce 和 Hive InputFormats 及 SerDes 实现(如需与 Hive 配合使用,请使用经过 shade 处理的 iceberg-hive-runtime 版本)
  • iceberg-nessie 是一个模块,用于将 Iceberg 表的元数据历史与操作集成到 Project Nessie 中
  • iceberg-data 是一个客户端库,用于从 JVM 应用程序读取 Iceberg 表
  • iceberg-runtime 会生成一个供 Spark 使用的 shade 运行时 jar,以便与 Iceberg 表集成

评论

登录后参与评论

正在加载评论…