探索Iceberg与Spark存储过程的深度整合之旅
Spark与Iceberg集成
表管理
Spark与Iceberg集成后,可以通过内置的存储过程来进行表的管理。使用CALL来调用存储过程。所有的存储过程在system的命名空间中。
由于表迁移功能的风险较大,所以不去进行表的迁移,使用重建Iceberg表,重写数据的方式进行切换。
调用语法
catalog_name代表catalog的名称procedure_name代表存储过程的名称- 参数可以通过指定参数名的方式入参,也可以使用位移的方式入参
sql
CALL catalog_name.system.procedure_name(arg_name_2 => arg_2, arg_name_1 => arg_1);
CALL catalog_name.system.procedure_name(arg_1, arg_2, ... arg_n);
调用样例
java
SparkSession spark = SparkSession
.builder()
.master("local")
.appName("Iceberg spark example")
.config("spark.sql.extensions", "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions")
.config("spark.sql.catalog.local", "org.apache.iceberg.spark.SparkCatalog")
.config("spark.sql.catalog.local.type", "hadoop") // 指定catalog类型
.config("spark.sql.catalog.local.warehouse", "iceberg_warehouse")
.getOrCreate();
spark.sql("CALL local.system.rollback_to_snapshot('iceberg_db.table2', 3285133177610707025)");
表快照管理
快照回滚
根据snapshotid进行回滚
存储过程名:rollback_to_snapshot

根据timestamp进行回滚
存储过程名:rollback_to_timestamp

设置表当前生效的快照
存储过程名:set_current_snapshot

❗️ snapshot_id与ref只能取其一。
表元数据管理
设置快照过期时间
Iceberg中的每次write/update/delete/upsert/compaction都会生成一个新快照,同时保留旧数据和元数据,以便进行快照隔离和时间旅行。expire_snapshots过程可用于删除不再需要的旧快照及其文件。
这个过程将删除旧快照和那些旧快照唯一需要的数据文件。这意味着expire_snapshots过程永远不会删除未过期快照仍然需要的文件。
存储过程名:expire_snapshots

如果省略older_than和retain_last,则将使用表的expiration properties。仍被分支或标记引用的快照不会被删除。默认情况下,分支和标记永不过期,但可以使用表属性history.expire.max-ref-age-ms更改其保留策略。main分支永不过期。
❗️ 使用此存储过程时,必须增加stream_results且值为true。
清除孤岛文件
用于删除未在Iceberg表的任何元数据文件中引用的文件,因此可视为“孤岛”。
存储过程名:remove_orphan_files

重写数据文件
Iceberg在一个表格中跟踪每个数据文件。数据文件越多,存储在清单文件中的元数据也就越多,而数据文件过小则会导致不必要的元数据量和文件打开成本,从而降低查询效率。
Iceberg可以使用Spark的rewriteDataFiles操作并行压缩数据文件。这将把小文件合并为大文件,以减少元数据开销和运行时文件打开成本。
存储过程名:rewrite_data_files

运用参数示例
sql
spark.sql("CALL catalog_name.system.rewrite_data_files(table => 'db.sample', options => map('min-input-files','2'))");
General Options

Options for sort strategy

Options for sort strategy with zorder sort_order

重写清单文件
重写表的清单,优化扫描规划。清单中的数据文件按分区规范中的字段排序。该程序使用Spark作业并行运行。
存储过程名:rewrite_manifests

重写位置删除文件
Iceberg可以重写位置删除文件,这样做有两个目的:
- 小型压缩:将小的位置删除文件压缩成大文件。这样可以减少存储在清单文件中的元数据大小,并减少打开小的删除文件的开销。
- 删除悬而未决的删除记录:过滤掉引用不再有效的数据文件的位置删除记录。重写数据文件后,指向重写数据文件的位置删除记录并不总是被标记为删除,而是会继续被表的实时快照元数据跟踪。这就是所谓的“悬空删除”问题。
存储过程名:rewrite_position_delete_files

在重写过程中,悬挂删除总是会被过滤掉。
Options

Iceberg v2表常见配置
表属性
读取属性

写入属性



表行为属性

保留表属性
保留表属性仅用于在创建或更新表时控制行为。这些属性的值不会作为表元数据的一部分进行持久化。

兼容性标志

Catalog 属性
IcebergCatalog支持使用Catalog属性来配置Catalog行为。下面是常用catalog属性的列表:

HadoopCatalog和HiveCatalog可以在其构造函数中访问这些属性。任何其他自定义Catalog可以通过实现Catalog.initialize(catalogName, catalogProperties)来访问这些属性。属性可以手动构建或从计算引擎(如Spark或Flink)传递。Spark使用其会话属性作为catalog属性,在Spark配置部分中可以找到更多详细信息。Flink通过CREATE CATALOG语句传递目录属性,有关更多详细信息,请参见Flink部分。
锁定Catalog 属性
以下是与锁定相关的Catalog属性。某些Catalog实现使用它们来控制提交期间的锁定行为。

Hadoop配置
Hive Metastore连接器使用Hadoop配置中的以下属性。HMS表锁定是一个两步过程:
- 锁创建:在HMS中创建锁并排队等待获取
- 锁检查:检查是否成功获取锁

注意: iceberg.hive.lock-check-max-wait-ms和iceberg.hive.lock-heartbeat-interval-ms应该小于Hive Metastore的事务超时时间(在较新的版本中为hive.txn.timeout或metastore.txn.timeout)。否则,在从Iceberg重新尝试锁之前,Hive Metastore中的锁心跳(发生在锁检查期间)会过期。
警告: 将iceberg.engine.hive.lock-enabled设置为false将导致HiveCatalog在提交表时不使用Hive锁。只有在满足以下所有条件时,才应将其设置为false:
- Hive Metastore服务器上可用
HIVE-26882 - 所有提交到与此
HiveCatalog提交的表相同的其他HiveCatalog都是Iceberg 1.3或更高版本 - 所有提交到与此
HiveCatalog提交的表相同的其他HiveCatalog在提交时也已禁用了Hive锁
未能确保这些条件会导致表损坏。即使将iceberg.engine.hive.lock-enabled设置为false,HiveCatalog仍然可以通过设置表属性engine.hive.lock-enabled=true来为单个表使用锁。这在其他HiveCatalog无法升级并设置为在提交时不使用Hive锁的情况下非常有用。
表Schema变更
Iceberg支持使用Alter table … alter column语法对Schema进行变更,示例如下:
sql
-- spark sql
-- 更改字段类型
ALTER TABLE prod.db.sample ALTER COLUMN measurement TYPE double;
-- 更新字段和 comment
ALTER TABLE prod.db.sample ALTER COLUMN measurement TYPE double COMMENT 'unit is bytes per second'
-- 更改字段顺序, FIRST/AFTER
ALTER TABLE prod.db.sample ALTER COLUMN col FIRST
ALTER TABLE prod.db.sample ALTER COLUMN nested.col AFTER other_col
-- null 更改,如果该字段是主键则不支持
ALTER TABLE prod.db.sample ALTER COLUMN id DROP NOT NULL
ALTER TABLE prod.db.sample ALTER COLUMN id SET NOT NULL
// 添加一个顶级列
table.UpdateSchema()
.addColumn("toplevel", Types.DecimalType.of(9, 2))
.commit()
// 添加一个列到嵌套类型里面
table.UpdateSchema()
.addColumn("points", "z", Types.LongType.get(), "z axis")
.commit()
// 更新列类型
table.UpdateSchema()
.updateColumn("id", Types.LongType.get())
.commit()
// 更新嵌套类里面子列类型
table.UpdateSchema()
.updateColumn("locations.lat", Types.DoubleType.get())
.commit()
table.UpdateSchema()
.updateColumn("locations.long", Types.DoubleType.get())
.commit()
// 重命名顶级列
table.UpdateSchema()
.renameColumn("data", "json")
.commit()
// 重命名嵌套类型子列
table.UpdateSchema()
.renameColumn("preferences.feature2", "newfeature")
.commit()
table.UpdateSchema()
.renameColumn("locations.lat", "latitude")
.commit()
table.UpdateSchema()
.renameColumn("points.x", "X")
.commit()
// 删除列
table.UpdateSchema()
.deleteColumn("points.z")
.commit();
Schema Visitor
java
// 从Spark表获取Schema并转成Iceberg Schema
Schema schema = SparkSchemaUtil
.schemaForTable(sparkSession, tableName);
// 从Spark表获取Partition信息,转成Iceberg PartitionSpec
PartitionSpec spc = SparkSchemaUtil
.specForTable(spark, tableName);
// Spark StructType和Iceberg Schema互相转换
Schema schema = SparkSchemaUtil.convert(structType);
StructType struct = SparkSchemaUtil.convert(schema);
Partition Evolution
scala
import org.apache.iceberg.expressions.Expressions
import org.apache.iceberg.Table
val table: Table = null
// =========Partition更新==========
table.updateSpec()
// 添加分区字段user_id,且一个分区内分桶数量为10
.addField(Expressions.bucket("user_id", 10))
// 删除分区字段country
.removeField("country")
// .commit()
// 可以通过 table = (BaseTable)origTable的方式获取内部table实例
TableMetadata current = table.operations().current();
PartitionSpec newSpec = PartitionSpec.builderFor(table.schema())
.hour("event_time")
.withSpecId(1)
.build();
table.ops().commit(current, current.updatePartitionSpec(newSpec));
Sort order evolution
java
Table sampleTable = ...;
sampleTable.replaceSortOrder()
.asc("id", NullOrder.NULLS_LAST)
.dec("category", NullOrder.NULL_FIRST)
.commit();
Iceberg 删除元数据中未被引用的文件(remove_orphan_files)
Iceberg使用JSON文件来跟踪表的元数据。每一次对表的更改都会生成一个新的元数据文件,以提供原子性。
默认情况下,旧的元数据文件会被保留以供历史记录。那些被流作业频繁提交的表可能需要定期清理元数据文件。要自动清理元数据文件,请在表属性中设置:
sql
write.metadata.delete-after-commit.enabled=true
write.metadata.previous-versions-max=5
write.metadata.delete-after-commit.enabled:是否在每次表提交后删除旧的跟踪元数据文件。write.metadata.previous-versions-max:保留的旧元数据文件的数量。
请注意,这只会删除在元数据日志中跟踪的元数据文件,不会删除孤立的元数据文件。需使用remove_orphan_files进行彻底删除。
在任何写入操作完成之前,使用比预期完成时间短的保留间隔来删除孤立文件是危险的,因为如果正在进行的文件被认为是孤立的并被删除,可能会破坏表。默认的间隔是3天。
Iceberg在确定哪些文件需要被移除时,使用路径的字符串表示形式。在一些文件系统上,路径随时间改变可能会变化,但它仍然代表同一个文件。例如,如果你更改了HDFS集群的权限,那么在创建期间使用的旧路径URL将不会与当前列表中出现的那些匹配。当运行RemoveOrphanFiles时,这将导致数据丢失。请确保你的MetadataTables中的条目与Hadoop FileSystem API列出的那些相匹配,以避免意外删除。
- 使用
expire_snapshots和remove_orphan_files清空旧快照并释放数据文件。 - 使用
rewrite_position_delete_files将小的位置删除文件合并为大文件。 - 使用
rewrite_data_files完全清除删除文件,从而实现最佳读取性能。
Spark 对Iceberg 删除元数据中未被引用的文件
自主方式合并小文件
执行顺序:rewrite_data_files -> rewrite_manifests -> expire_snapshots -> remove_orphan_files
sql
call spark_catalog.system.remove_orphan_files(table => 'db.sample');
-- 删除孤立文件
table.refresh()
val orphanFilesRemoved = table.optimize("RemoveOrphanFiles")
println(s"Orphan files removed: $orphanFilesRemoved")
合并数据文件
sql
spark.sql(s"""
|CALL $catalogName.system.rewrite_data_files('${metaTable.schema}.${metaTable.tableName}', 'sort', '$specName DESC NULLS LAST',
| map('min-input-files', '1', 'target-file-size-bytes', '${128 * 1024 * 1024}', 'partial-progress.enabled', 'true', 'partial-progress.max-commits', '100'),
| '$specName >= "${dateFormat.format(tsToExpire)}" and $specName <= "${dateFormat.format(currentTimeMillis)}"')
|""".stripMargin)
合并目录文件
sql
spark.sql(
s"""
|CALL $catalogName.system.rewrite_manifests('${metaTable.schema}.${metaTable.tableName}', false)
|""".stripMargin)
删除过期快照(包含data file和manifest file)
sql
spark.sql(
s"""
|CALL $catalogName.system.expire_snapshots('${metaTable.schema}.${metaTable.tableName}',
| TIMESTAMP '${dateFormat.format(tsToExpire)}', 12, 10)
|""".stripMargin)
删除孤立文件(不必每次都执行)
sql
spark.sql(
s"""
|CALL $catalogName.system.remove_orphan_files('${metaTable.schema}.${metaTable.tableName}',
| TIMESTAMP '${dateFormat.format(tsToExpire)}', '${table.location()}/data', false, 10)
|""".stripMargin)
通过表属性设置合并小文件
通过建表时初始化表属性TBLPROPERTIES ('write.distribution-mode'='hash'),可动态合并小文件
sql
spark.sql(
s"""
|CREATE TABLE IF NOT EXISTS hadoop_prod.$schemaName.$tableName (
| id bigint,
| data string,
| jydw_no string,
| ts timestamp)
|USING iceberg
|PARTITIONED BY (bucket(5, id), days(ts))
|TBLPROPERTIES ('write.distribution-mode'='hash',
|'write.metadata.delete-after-commit.enabled'='true',
|'write.metadata.previous-versions-max'='9')
|""".stripMargin)
Trino 删除元数据中未被引用的文件(remove_orphan_files)
sql
-- 清理快照
ALTER TABLE test_table EXECUTE remove_orphan_files(retention_threshold => '7d')
-- 合并小文件
ALTER TABLE test_table EXECUTE optimize(file_size_threshold => '10MB')
Flink 删除元数据中未被引用的文件(remove_orphan_files)
java
// 1. 定期合并新增分区的小文件:
rewriteDataFilesAction.execute(); 仅合并小文件,不会删除旧文件。
// 2. 删除过期的 snapshot,清理元数据及数据文件:
table.expireSnapshots().expireOlderThan(timestamp).commit();
// 3. 清理 orphan 文件,默认清理 3 天前,且无法触及的文件:
removeOrphanFilesAction.olderThan(timestamp).execute();
int maxSize = getMaxParallelSize();
System.out.println("start to removeOrphanFile before " + today + ", " + end);
ActionsWithParallel.forTable(table)
.removeOrphanFilesWithParallel()
.olderThan(end) // 谨慎操作,避免误删当前写入还未成功commit的datafile
.execute(maxSize, 10000);
// tablelocation 改为具体的表地址,会清理不用的文件,
CALL hive_iceberg_catalog.system.remove_orphan_files(table => 'ods_base.xx_table_name', location => 'tablelocation/data')
大数据相关学习资料
大数据相关学习资料、大数据项目、湖仓一体、架构师必知必会、数据中台建设方法论...共有1400多份文档资料,另专为星球成员整理了一份比较详细的语雀知识库合计190万字和飞书文档资料。(内容太多,此处不做展示...),欢迎大家踊跃加入星球,您将获得:
- 提供最全的大数据知识库,不限设备,随时随地打开看的在线文档。
- 免费答疑解惑、交流技术
- 面试指导、模拟面试
- 各类pdf文档下载、星球代码下载
- 提供简历模板,简历修改指导服务,星球成员免费提供简历修改指导。
另外说明加入星球后支持三天无理由退款,不满意无条件随时退。需要资料请加微信:D1435221412
end
