如何读取Apache Iceberg最新快照中的新增记录?
仅读取Apache Iceberg最新快照中新增记录的方法
你可以通过读取相邻快照的差异数据并过滤新增类型的记录来实现需求,以下是具体步骤:
1. 获取快照ID对
首先需要拿到最新快照ID及其父快照ID(即上一个版本的快照ID),可以通过查询Iceberg的快照元数据表获取:
SELECT snapshot_id, parent_id FROM your_table.snapshots ORDER BY committed_at DESC LIMIT 2;
这条SQL会返回最新的两个快照,第一条是当前最新快照,第二条的snapshot_id就是父快照ID。
2. 读取快照差异并过滤新增记录
使用Spark读取这两个快照之间的差异数据,然后通过Iceberg自动生成的_change_type字段筛选出新增的记录:
spark.read .format("iceberg") .option("start-snapshot-id", "父快照ID") .option("end-snapshot-id", "最新快照ID") .load("path/to/table") .filter("_change_type = 'insert'")
_change_type的可选值说明:
insert:纯新增的记录delete:被删除的记录update_before:更新前的旧记录update_after:更新后的新记录
如果你的场景需要包含更新后的最新行(而非仅纯插入),可以把过滤条件改成_change_type IN ('insert', 'update_after')。
3. 自动获取快照ID(代码层面)
如果需要在代码中自动获取快照ID,避免手动查询,可以借助Iceberg的Scala API:
import org.apache.iceberg.spark.SparkCatalog import org.apache.iceberg.TableIdentifier // 加载Iceberg表 val table = spark.sessionState.catalog.asInstanceOf[SparkCatalog] .loadTable(TableIdentifier.of("你的数据库名", "你的表名")) // 获取最新快照和父快照ID val latestSnapshot = table.currentSnapshot() val parentSnapshotId = latestSnapshot.parentId() // 读取新增记录 val newRecords = spark.read .format("iceberg") .option("start-snapshot-id", parentSnapshotId.toString) .option("end-snapshot-id", latestSnapshot.snapshotId().toString) .load("你的数据库名.你的表名") .filter("_change_type = 'insert'")
关于你之前方法的说明
你之前用start-snapshot-id和end-snapshot-id的方法本身支持非追加模式,但它会返回两个快照之间的所有变更记录(包括删除、更新的前后行),所以需要通过_change_type过滤才能得到纯新增的部分,并非仅适用于追加模式。
内容的提问来源于stack exchange,提问作者karas27
相关产品推荐
相关产品推荐

