API
Iceberg Java API
表
Iceberg API 的主要用途是管理表元数据,例如模式(schema)、分区规范、元数据以及存储表数据的数据文件。
表元数据和相关操作通过 Table 接口访问,该接口用于返回表信息。
表元数据
Table 接口提供对表元数据的访问:
schema返回当前表的模式spec返回当前表的分区规范properties返回键值对形式的属性映射currentSnapshot返回当前表的快照snapshots返回该表的所有有效快照snapshot(id)根据 ID 返回指定快照location返回表的根目录位置
表还提供 refresh 方法,用于将表更新到最新版本,并暴露以下辅助方法:
io返回用于读写表文件的FileIOlocationProvider返回用于生成数据文件和元数据文件路径的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 中的工厂方法。
支持的谓词表达式有:
isNullnotNullequalnotEquallessThanlessThanOrEqualgreaterThangreaterThanOrEqualinnotInstartsWithnotStartsWith
支持的表达式运算有:
andornot
常量表达式有:
alwaysTruealwaysFalse
表达式绑定
表达式在创建时处于未绑定状态。在使用表达式之前,需要将其绑定到某个数据类型,以查找表达式名称所代表的字段 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 表集成
评论
登录后参与评论
KnowForge