快速入门
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);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()评论
登录后参与评论
KnowForge