API

快速入门

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

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

Java API 快速入门

创建表

表可以通过 Catalog 或 Tables 接口的实现来创建。

使用 Hive catalog

Hive catalog 连接到 Hive metastore,用于跟踪 Iceberg 表。你可以使用一个名称和若干属性来初始化 Hive catalog。(参见:Catalog 属性)

import java.util.HashMap;
import java.util.Map;

import org.apache.iceberg.hive.HiveCatalog;

HiveCatalog catalog = new HiveCatalog();
catalog.setConf(spark.sparkContext().hadoopConfiguration());  // Optionally use Spark's Hadoop configuration

Map <String, String> properties = new HashMap<String, String>();
properties.put("warehouse", "...");
properties.put("uri", "...");

catalog.initialize("hive", properties);

HiveCatalog 实现了 Catalog 接口,该接口定义了用于操作表的方法,例如 createTable、loadTable、renameTable 和 dropTable。要创建一张表,需要传入一个 Identifier 和一个 Schema,以及其他初始元数据:

import org.apache.iceberg.Table;
import org.apache.iceberg.catalog.TableIdentifier;

TableIdentifier name = TableIdentifier.of("logging", "logs");
Table table = catalog.createTable(name, schema, spec);

// or to load an existing table, use the following line
Table table = catalog.loadTable(name);

该表的 schema 和 partition spec 在下方创建。

使用 Hadoop catalog

Hadoop catalog 不需要连接 Hive MetaStore,但只能与 HDFS 或其他支持原子重命名的文件系统一起使用。使用本地文件系统或 S3 时,通过 Hadoop catalog 进行并发写入是不安全的。要创建 Hadoop catalog:

import org.apache.hadoop.conf.Configuration;
import org.apache.iceberg.hadoop.HadoopCatalog;

Configuration conf = new Configuration();
String warehousePath = "hdfs://host:8020/warehouse_path";
HadoopCatalog catalog = new HadoopCatalog(conf, warehousePath);

与 Hive catalog 一样,HadoopCatalog 也实现了 Catalog 接口,因此它同样提供了操作表的方法,例如 createTable、loadTable 和 dropTable。

下面的示例使用 Hadoop catalog 创建一张表:

import org.apache.iceberg.Table;
import org.apache.iceberg.catalog.TableIdentifier;

TableIdentifier name = TableIdentifier.of("logging", "logs");
Table table = catalog.createTable(name, schema, spec);

// or to load an existing table, use the following line
Table table = catalog.loadTable(name);

表的 schema 和 分区规范 在下方创建。

Spark 中的表

Spark 可以通过 HiveCatalog 按名称操作表。

// spark.sql.catalog.hive_prod = org.apache.iceberg.spark.SparkCatalog
// spark.sql.catalog.hive_prod.type = hive
spark.table("logging.logs");

Spark 还可以通过路径加载 HadoopCatalog 创建的表。

spark.read.format("iceberg").load("hdfs://host:8020/warehouse_path/logging/logs");

Schema

创建 Schema

下面的示例为 logs 表创建一个 Schema:

import org.apache.iceberg.Schema;
import org.apache.iceberg.types.Types;

Schema schema = new Schema(
      Types.NestedField.required(1, "level", Types.StringType.get()),
      Types.NestedField.required(2, "event_time", Types.TimestampType.withZone()),
      Types.NestedField.required(3, "message", Types.StringType.get()),
      Types.NestedField.optional(4, "call_stack", Types.ListType.ofRequired(5, Types.StringType.get()))
    );

直接使用 Iceberg API 时,必须提供类型 ID。从其他 schema 格式(如 Spark、Avro 和 Parquet)转换时,会自动分配新的 ID。

创建表时,schema 中的所有 ID 都会重新分配,以确保唯一性。

从 Avro 转换 schema

要从现有的 Avro schema 创建 Iceberg schema,请使用 AvroSchemaUtil 中的转换器:

import org.apache.avro.Schema;
import org.apache.avro.Schema.Parser;
import org.apache.iceberg.avro.AvroSchemaUtil;

Schema avroSchema = new Parser().parse("{\"type\": \"record\" , ... }");
Schema icebergSchema = AvroSchemaUtil.toIceberg(avroSchema);

从 Spark 转换 schema

要从现有表创建 Iceberg schema,请使用 SparkSchemaUtil 中的转换器:

import org.apache.iceberg.spark.SparkSchemaUtil;

Schema schema = SparkSchemaUtil.schemaForTable(sparkSession, tableName);

分区

创建分区规格

分区规格(Partition spec)描述 Iceberg 应如何将记录分组到数据文件中。分区规格针对表的架构,通过构建器(builder)创建。

以下示例为 logs 表创建分区规格,按日志事件时间戳的小时以及日志级别对记录进行分区:

import org.apache.iceberg.PartitionSpec;

PartitionSpec spec = PartitionSpec.builderFor(schema)
      .hour("event_time")
      .identity("level")
      .build();

欲了解 Iceberg 提供的各种分区转换的更多信息,请访问此页面。

分支与标签

创建分支和标签

可以通过 Java 库的 ManageSnapshots API 创建新的分支和标签。

/* Create a branch test-branch which is retained for 1 week, and the latest 2 snapshots on test-branch will always be retained.
Snapshots on test-branch which are created within the last hour will also be retained. */

String branch = "test-branch";
table.manageSnapshots()
    .createBranch(branch, 3)
    .setMinSnapshotsToKeep(branch, 2)
    .setMaxSnapshotAgeMs(branch, 3600000)
    .setMaxRefAgeMs(branch, 604800000)
    .commit();

// Create a tag historical-tag at snapshot 10 which is retained for a day
String tag = "historical-tag"
table.manageSnapshots()
    .createTag(tag, 10)
    .setMaxRefAgeMs(tag, 86400000)
    .commit();

提交到分支

通过在操作中指定 toBranch 即可向分支写入。完整列表请参阅 UpdateOperations。

// Append FILE_A to branch test-branch
String branch = "test-branch";

table.newAppend()
    .appendFile(FILE_A)
    .toBranch(branch)
    .commit();


// Perform row level updates on "test-branch"
table.newRowDelta()
    .addRows(DATA_FILE)
    .addDeletes(DELETES)
    .toBranch(branch)
    .commit();


// Perform a rewrite operation replacing SMALL_FILE_1 and SMALL_FILE_2 on "test-branch" with compactedFile.
table.newRewrite()
    .rewriteFiles(ImmutableSet.of(SMALL_FILE_1, SMALL_FILE_2), ImmutableSet.of(compactedFile))
    .toBranch(branch)
    .commit();

从分支和标签读取

从分支或标签读取的方式与常规相同,通过 Table Scan API 即可实现,只需在 useRef API 中传入分支或标签即可。当传入分支时,使用的快照是该分支的当前最新快照(head)。请注意,目前尚不支持在从分支读取的同时在扫描中指定 asOfSnapshotId。

// Read from the head snapshot of test-branch
TableScan branchRead = table.newScan().useRef("test-branch");

// Read from the snapshot referenced by audit-tag
TableScan tagRead = table.newScan().useRef("audit-tag");

替换与快进分支和标签

现有分支和标签所指向的快照可以通过 replace API 进行更新。快进操作与 git 的快进类似。当目标分支是源分支的祖先时,可以使用快进将目标分支推进到源分支或标签的最新提交。对于快进和替换操作,默认都会保留目标分支的保留属性。

// Update "test-branch" to point to snapshot 4
table.manageSnapshots()
     .replaceBranch(branch, 4)
     .commit()

String tag = "audit-tag";
// Replace "audit-tag" to point to snapshot 3 and update its retention
table.manageSnapshots()
     .replaceBranch(tag, 4)
     .setMaxRefAgeMs(1000)
     .commit()

更新保留策略属性

分支和标签的保留策略属性同样可以更新。使用 setMaxRefAgeMs 来更新分支或标签自身的保留属性。分支的快照保留属性可以通过 setMinSnapshotsToKeep 和 setMaxSnapshotAgeMs API 进行更新。

String branch = "test-branch";
// Update retention properties for test-branch
table.manageSnapshots()
    .setMinSnapshotsToKeep(branch, 10)
    .setMaxSnapshotAgeMs(branch, 7200000)
    .setMaxRefAgeMs(branch, 604800000)
    .commit();

// Update retention properties for test-tag
table.manageSnapshots()
    .setMaxRefAgeMs("test-tag", 604800000)
    .commit();

删除分支与标签

分支和标签可以通过 removeBranch 和 removeTag API 分别删除。

// Remove test-branch
table.manageSnapshots()
     .removeBranch("test-branch")
     .commit()

// Remove test-tag
table.manageSnapshots()
     .removeTag("test-tag")
     .commit()

评论

登录后参与评论

正在加载评论…