改版通知

巨人肩膀网站已全新改版。若您仍依赖旧站功能或数据,欢迎联系我们,我们会协助处理。联系我们

探索Iceberg与Spark存储过程的深度整合之旅

ckckck2025年1月10日52 浏览

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_idref只能取其一。

表元数据管理

设置快照过期时间

Iceberg中的每次write/update/delete/upsert/compaction都会生成一个新快照,同时保留旧数据和元数据,以便进行快照隔离和时间旅行。expire_snapshots过程可用于删除不再需要的旧快照及其文件。

这个过程将删除旧快照和那些旧快照唯一需要的数据文件。这意味着expire_snapshots过程永远不会删除未过期快照仍然需要的文件。

存储过程名:expire_snapshots

如果省略older_thanretain_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属性的列表:

HadoopCatalogHiveCatalog可以在其构造函数中访问这些属性。任何其他自定义Catalog可以通过实现Catalog.initialize(catalogName, catalogProperties)来访问这些属性。属性可以手动构建或从计算引擎(如Spark或Flink)传递。Spark使用其会话属性作为catalog属性,在Spark配置部分中可以找到更多详细信息。Flink通过CREATE CATALOG语句传递目录属性,有关更多详细信息,请参见Flink部分。

锁定Catalog 属性

以下是与锁定相关的Catalog属性。某些Catalog实现使用它们来控制提交期间的锁定行为。

Hadoop配置

Hive Metastore连接器使用Hadoop配置中的以下属性。HMS表锁定是一个两步过程:

  1. 锁创建:在HMS中创建锁并排队等待获取
  2. 锁检查:检查是否成功获取锁

注意: iceberg.hive.lock-check-max-wait-msiceberg.hive.lock-heartbeat-interval-ms应该小于Hive Metastore的事务超时时间(在较新的版本中为hive.txn.timeoutmetastore.txn.timeout)。否则,在从Iceberg重新尝试锁之前,Hive Metastore中的锁心跳(发生在锁检查期间)会过期。

警告:iceberg.engine.hive.lock-enabled设置为false将导致HiveCatalog在提交表时不使用Hive锁。只有在满足以下所有条件时,才应将其设置为false

  1. Hive Metastore服务器上可用HIVE-26882
  2. 所有提交到与此HiveCatalog提交的表相同的其他HiveCatalog都是Iceberg 1.3或更高版本
  3. 所有提交到与此HiveCatalog提交的表相同的其他HiveCatalog在提交时也已禁用了Hive锁

未能确保这些条件会导致表损坏。即使将iceberg.engine.hive.lock-enabled设置为falseHiveCatalog仍然可以通过设置表属性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列出的那些相匹配,以避免意外删除。

  1. 使用expire_snapshotsremove_orphan_files清空旧快照并释放数据文件。
  2. 使用rewrite_position_delete_files将小的位置删除文件合并为大文件。
  3. 使用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')  
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万字和飞书文档资料。(内容太多,此处不做展示...),欢迎大家踊跃加入星球,您将获得:

  1. 提供最全的大数据知识库,不限设备,随时随地打开看的在线文档。
  2. 免费答疑解惑、交流技术
  3. 面试指导、模拟面试
  4. 各类pdf文档下载、星球代码下载
  5. 提供简历模板,简历修改指导服务,星球成员免费提供简历修改指导。

另外说明加入星球后支持三天无理由退款,不满意无条件随时退。需要资料请加微信:D1435221412
end