拓十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Apache Iceberg Spark 查询指南:SQL、DataFrame、时态旅行与元数据表实战

Apache Iceberg Spark 查询指南:SQL、DataFrame、时态旅行与元数据表实战 数据湖大数据数据存储【免费下载链接】icebergApache Iceberg项目地址https://gitcode.com/gh_mirrors/icebe/iceberg点击查看免费下载Apache Iceberg 通过 Spark 的 DataSourceV2 API 实现了完整的查询能力本文基于仓库中的 spark-queries.md 文档系统讲解在 Spark 上查询 Iceberg 表的四种核心方式标准 SQL 查询、Iceberg 特有的 SQL 变换函数、基于 DataFrame API 的读取以及用于表诊断的元数据表体系。读完本文你将掌握如何用 SQL 和 DataFrame 对 Iceberg 表做时态旅行time travel、增量读取、按分支/标签读取以及如何借助history、snapshots、files、manifests、partitions等元数据表深入排查表状态并理解其背后的源码实现。前置准备配置 Spark Catalog要在 Spark 中使用 Iceberg第一步是配置 Spark catalogs。Iceberg 使用 Apache Spark 的 DataSourceV2 API 来实现数据源data source与目录catalog实现。查询时表标识符identifier包含目录名形如prod.db.table目录prod、命名空间db、表tableSELECT * FROM prod.db.table; -- catalog: prod, namespace: db, table: tableIceberg 提供两类目录实现见 spark-configuration.mdorg.apache.iceberg.spark.SparkCatalog支持 Hive、Hadoop、REST、Glue、JDBC、Nessie 等目录类型org.apache.iceberg.spark.SparkSessionCatalog为 Spark 内置 catalogspark_catalog添加 Iceberg 表支持非 Iceberg 表仍委托给内置 catalog 处理。例如配置一个 Hive 目录spark.sql.catalog.hive_prod org.apache.iceberg.spark.SparkCatalog spark.sql.catalog.hive_prod.type hive spark.sql.catalog.hive_prod.uri thrift://metastore-host:port # 省略 uri 时使用 hive-site.xml 中与 Spark 相同的 URIhive.metastore.uris使用 SQL 查询元数据表的 SQL 访问Iceberg 的元数据表metadata tables可以像普通表一样通过 SQL 读取——把元数据表名作为 Iceberg 表名的一个命名空间后缀。例如读取prod.db.table的files元数据表SELECT * FROM prod.db.table.files;该查询返回每个数据文件的行级元数据包括文件路径、格式、分区、记录数、文件大小、列统计信息column_sizes、value_counts、null_value_counts、lower_bounds、upper_bounds等。这也是后续检查表章节的基础。Iceberg SQL 变换函数Iceberg 为每个 Iceberg catalog 添加了一组 SQL 函数用于在查询中检查变换transform结果以及编写与 Iceberg 分区变换匹配的过滤器。这些函数只通过 Iceberg catalog 可用不会注册到 Spark 内置 catalog。注意Spark 4.2.0 之前的版本不支持 session catalog 中的V2Function。即使将spark_catalog配置为org.apache.iceberg.spark.SparkSessionCatalogSELECT spark_catalog.system.bucket(16, id)这类查询仍会失败详见 SPARK-54760 与 apache/spark#53531。要使用 Iceberg SQL 函数请通过配置了org.apache.iceberg.spark.SparkCatalog的 catalog 调用。调用这些函数时使用system命名空间SELECT system.iceberg_version(); SELECT system.bucket(16, id), system.days(ts) FROM prod.db.table;想显式指明目录时可以在函数前加上目录名限定SELECT prod.system.bucket(16, id) FROM prod.db.table;提示PARTITIONED BY子句使用单数变换表达式如year(ts)、month(ts)而 SQL 函数使用system.years(ts)、system.months(ts)。完整函数清单如下函数支持的输入类型返回类型示例system.iceberg_version()无stringSELECT system.iceberg_version();system.bucket(numBuckets, col)date、tinyint、smallint、int、bigint、timestamp、timestamp_ntz、decimal、string、binaryintSELECT system.bucket(16, id) FROM prod.db.table;system.years(col)date、timestamp、timestamp_ntzintSELECT system.years(ts) FROM prod.db.table;system.months(col)date、timestamp、timestamp_ntzintSELECT system.months(ts) FROM prod.db.table;system.days(col)date、timestamp、timestamp_ntzdateSELECT * FROM prod.db.table WHERE system.days(ts) date(2025-03-01);system.hours(col)timestamp、timestamp_ntzintSELECT system.hours(ts) FROM prod.db.table;system.truncate(width, col)tinyint、smallint、int、bigint、decimal、string、binary与col相同类型SELECT system.truncate(4, data) FROM prod.db.table;所有变换函数对NULL输入都返回NULL。语义要点system.years、system.months、system.days、system.hours返回的是Iceberg 变换值而非日历字段。例如system.years返回自 1970-01-01 起的年数system.months返回自 1970-01 起的月数system.hours返回自 1970-01-01T00:00 起的小时数。system.days返回表示输入日期部分的date值对date输入原样返回对 timestamp 输入丢弃时间分量。对数值输入system.truncate(width, col)向下取整到width的最近倍数对string和binary输入保留前width个字符或字节。这些函数在两种场景下尤其有用检查 Iceberg 如何变换某个值以及编写与分区变换对齐的查询和行级操作过滤器。从源码看这些函数在 SparkFunctions.java 中注册为years、months、days、hours、bucket、truncate、iceberg_version七个UnboundFunction通过 SparkFunctionCatalog.java 挂载到 Iceberg catalog函数解析不区分大小写load时统一toLowerCase。对应的实现类如 BucketFunction.java、TruncateFunction.java在 TestSparkFunctions.java 中被逐一验证覆盖了 date/int/long/timestamp/decimal/string/binary 等各类型变体。SQL 时态旅行Time TravelSpark 在 SQL 查询中通过TIMESTAMP AS OF或VERSION AS OF子句支持时态旅行。VERSION AS OF子句可以包含 long 类型的快照 ID或字符串形式的分支branch/标签tag名。注意如果分支或标签的名字与某个快照 ID 相同时态旅行会选择给定快照 ID 对应的快照。例如存在名为1的标签指向快照 ID 2那么VERSION AS OF 1会旅行到快照 ID 为 1 的快照。若不符合预期请用snapshot-1这类带明确前缀的名字重命名标签或分支。-- 时态旅行到 1986-10-26 01:21:00 SELECT * FROM prod.db.table TIMESTAMP AS OF 1986-10-26 01:21:00; -- 时态旅行到快照 ID 为 10963874102873L 的快照 SELECT * FROM prod.db.table VERSION AS OF 10963874102873; -- 时态旅行到 audit-branch 分支的头部快照 SELECT * FROM prod.db.table VERSION AS OF audit-branch; -- 时态旅行到标签 historical-snapshot 引用的快照 SELECT * FROM prod.db.table VERSION AS OF historical-snapshot;此外FOR SYSTEM_TIME AS OF与FOR SYSTEM_VERSION AS OF子句同样受支持SELECT * FROM prod.db.table FOR SYSTEM_TIME AS OF 1986-10-26 01:21:00; SELECT * FROM prod.db.table FOR SYSTEM_VERSION AS OF 10963874102873; SELECT * FROM prod.db.table FOR SYSTEM_VERSION AS OF audit-branch; SELECT * FROM prod.db.table FOR SYSTEM_VERSION AS OF historical-snapshot;时间戳也可以用 Unix 时间戳秒提供-- 以秒为单位的时间戳 SELECT * FROM prod.db.table TIMESTAMP AS OF 499162860; SELECT * FROM prod.db.table FOR SYSTEM_TIME AS OF 499162860;分支或标签还可以用类似元数据表的语法指定即branch_branchname或tag_tagnameSELECT * FROM prod.db.table.branch_audit-branch; SELECT * FROM prod.db.table.tag_historical-snapshot;含-的标识符不合法因此必须用反引号转义。注意带分支/标签的标识符不能与VERSION AS OF组合使用。时态旅行查询中的 schema 选择不同类型的时态旅行查询会使用快照 schema 或表 schema-- 时态旅行到 1986-10-26 01:21:00 - 使用快照的 schema SELECT * FROM prod.db.table TIMESTAMP AS OF 1986-10-26 01:21:00; -- 时态旅行到快照 ID 为 10963874102873L 的快照 - 使用快照的 schema SELECT * FROM prod.db.table VERSION AS OF 10963874102873; -- 时态旅行到 audit-branch 分支头部 - 使用表的 schema SELECT * FROM prod.db.table VERSION AS OF audit-branch; SELECT * FROM prod.db.table.branch_audit-branch; -- 时态旅行到标签 historical-snapshot 引用的快照 - 使用快照的 schema SELECT * FROM prod.db.table VERSION AS OF historical-snapshot; SELECT * FROM prod.db.table.tag_historical-snapshot;用一个随时间演化的表来验证每种时态旅行查询的 schema 选择行为-- 快照 S1初始 schema (id, status) CREATE TABLE prod.db.orders ( id BIGINT, status STRING ) USING iceberg; INSERT INTO prod.db.orders VALUES (1, NEW), (2, PAID); -- 记录 S1 的 snapshot_id 与 committed_at -- 例如 snapshot_id 101, committed_at 2025-01-01 10:00:00 -- 快照 S2新增列 total 并写入新数据 ALTER TABLE prod.db.orders ADD COLUMN total DOUBLE; INSERT INTO prod.db.orders VALUES (3, PAID, 100.0); -- 此时 S2 为当前快照schema 为 (id, status, total)选择特定快照或时间戳的时态旅行查询使用快照的 schema-- 使用 S1 的快照 schema列为 (id, status) SELECT * FROM prod.db.orders VERSION AS OF 101; SELECT * FROM prod.db.orders TIMESTAMP AS OF 2025-01-01 10:00:00;上面两个查询的结果都只有id和status。total列在 S1 的 schema 中不存在即使当前表 schema 已包含total该列也不可见。现在创建一个分支和一个标签二者都指向 S1-- 分支 audit_branch 指向快照 S1 ALTER TABLE prod.db.orders CREATE BRANCH audit_branch AS OF VERSION 101; -- 标签 first_load 也指向快照 S1 ALTER TABLE prod.db.orders CREATE TAG first_load AS OF VERSION 101;查询分支时Spark 使用表的当前 schema-- 使用表 schema列为 (id, status, total) SELECT * FROM prod.db.orders VERSION AS OF audit_branch; -- 等价的标识符写法 SELECT * FROM prod.db.orders.branch_audit_branch;这些查询的结果包含(id, status, total)三列。对于来自 S1 的行total返回NULL因为写入这些行时该列尚不存在。查询标签时Spark 使用标签所引用快照的 schema-- 使用 S1 的快照 schema列为 (id, status) SELECT * FROM prod.db.orders VERSION AS OF first_load; -- 等价的标识符写法 SELECT * FROM prod.db.orders.tag_first_load;这些查询只返回id和status因为标签绑定到特定快照并使用该快照的 schema即使表当前 schema 已经演化。从源码看SQL 层面对时态旅行选项的解析最终汇入 SparkScanBuilder.java 的buildBatchScansnapshot-id对应scan.useSnapshot(snapshotId)as-of-timestamp对应scan.asOfTime(asOfTimestamp)branch/tag对应scan.useRef(branch)或scan.useRef(tag)当同时设置snapshot-id与as-of-timestamp时会抛出 Cannot set both ... to select which table snapshot to scan 错误。使用 DataFrame 查询要把表加载为 DataFrame使用table方法val df spark.table(prod.db.table)DataFrameReader 加载目录路径和表名可以通过 Spark 的DataFrameReader接口加载。使用spark.read.format(iceberg).load(table)或spark.table(table)时table变量可取下列形式按优先级从高到低排列例如匹配的 catalog 优先于任何命名空间解析file:///path/to/table加载指定路径下的 HadoopTabletablename加载currentCatalog.currentNamespace.tablenamecatalog.tablename从指定 catalog 加载tablenamenamespace.tablename从当前 catalog 加载namespace.tablenamecatalog.namespace.tablename从指定 catalog 加载namespace.tablenamenamespace1.namespace2.tablename从当前 catalog 加载namespace1.namespace2.tablenameDataFrame 时态旅行要在 DataFrame API 中选择特定表快照或某个时间点的快照Iceberg 支持四个 Spark 读取选项snapshot-id选择特定的表快照as-of-timestamp选择指定时间戳毫秒时的当前快照branch选择指定分支的头部快照。注意目前branch不能与as-of-timestamp组合使用tag选择与指定标签关联的快照。标签不能与as-of-timestamp组合使用// 时态旅行到 1986-10-26 01:21:00 spark.read .option(as-of-timestamp, 499162860000) .format(iceberg) .load(path/to/table)// 时态旅行到快照 ID 为 10963874102873L 的快照 spark.read .option(snapshot-id, 10963874102873L) .format(iceberg) .load(path/to/table)// 时态旅行到标签 historical-snapshot spark.read .option(SparkReadOptions.TAG, historical-snapshot) .format(iceberg) .load(path/to/table)// 时态旅行到 audit-branch 分支的头部快照 spark.read .option(SparkReadOptions.BRANCH, audit-branch) .format(iceberg) .load(path/to/table)这些选项在 SparkReadOptions.java 中定义为常量SNAPSHOT_ID snapshot-id、AS_OF_TIMESTAMP as-of-timestamp、BRANCH branch、TAG tag。增量读取Incremental read增量读取追加的数据使用以下选项start-snapshot-id增量扫描的起始快照 ID不含该快照本身end-snapshot-id增量扫描的结束快照 ID包含该快照。可选省略时默认使用当前快照// 读取 start-snapshot-id (10963874102873L) 之后、end-snapshot-id (63874143573109L) 之前新增的数据 spark.read .format(iceberg) .option(start-snapshot-id, 10963874102873) .option(end-snapshot-id, 63874143573109) .load(path/to/table)注意增量读取目前只能获取append操作产生的数据不支持replace、overwrite、delete操作。增量读取对 V1 和 V2 format-version 均适用但Spark 的 SQL 语法不支持增量读取。从源码看增量读取的约束在 TestDataSourceOptions.java 中有明确验证同时设置start-snapshot-id与snapshot-id/as-of-timestamp会报 Cannot set start-snapshot-id and end-snapshot-id for incremental scans when either snapshot-id or as-of-timestamp is set只设置end-snapshot-id而缺少start-snapshot-id会报 Cannot set only end-snapshot-id for incremental scans. Please, set start-snapshot-id too.。检查表元数据表全览要检查表的历史、快照及其他元数据Iceberg 支持元数据表metadata tables。元数据表通过在原始表名后追加元数据表名来标识例如db.table的历史用db.table.history读取。History表历史SELECT * FROM prod.db.table.history;返回made_current_at、snapshot_id、parent_id、is_current_ancestor等列。示例输出| made_current_at | snapshot_id | parent_id | is_current_ancestor | | -- | -- | -- | -- | | 2019-02-08 03:29:51.215 | 5781947118336215154 | NULL | true | | 2019-02-08 03:47:55.948 | 5179299526185056830 | 5781947118336215154 | true | | 2019-02-09 16:24:30.13 | 296410040247533544 | 5179299526185056830 | false | | 2019-02-09 16:32:47.336 | 2999875608062437330 | 5179299526185056830 | true | | 2019-02-09 19:42:03.919 | 8924558786060583479 | 2999875608062437330 | true | | 2019-02-09 19:49:16.343 | 6536733823181975045 | 8924558786060583479 | true |提示上表展示了一次被回滚的提交。示例中有两个快照拥有同一个父快照其中一个是不是当前表状态的祖先。Metadata Log Entries元数据日志条目SELECT * from prod.db.table.metadata_log_entries;返回timestamp、file、latest_snapshot_id、latest_schema_id、latest_sequence_number等列记录每次元数据文件.metadata.json的更新| timestamp | file | latest_snapshot_id | latest_schema_id | latest_sequence_number | | -- | -- | -- | -- | -- | | 2022-07-28 10:43:52.93 | s3://.../table/metadata/00000-9441e604-b3c2-498a-a45a-6320e8ab9006.metadata.json | null | null | null | | 2022-07-28 10:43:57.487 | s3://.../table/metadata/00001-f30823df-b745-4a0a-b293-7532e0c99986.metadata.json | 170260833677645300 | 0 | 1 | | 2022-07-28 10:43:58.25 | s3://.../table/metadata/00002-2cc2837a-02dc-4687-acc1-b4d86ea486f4.metadata.json | 958906493976709774 | 0 | 2 |Snapshots快照SELECT * FROM prod.db.table.snapshots;返回committed_at、snapshot_id、parent_id、operation、manifest_list、summary等列。示例| committed_at | snapshot_id | parent_id | operation | manifest_list | summary | | -- | -- | -- | -- | -- | -- | | 2019-02-08 03:29:51.215 | 57897183625154 | null | append | s3://.../table/metadata/snap-57897183625154-1.avro | { added-records - 2478404, total-records - 2478404, added-data-files - 438, total-data-files - 438, spark.app.id - application_1520379288616_155055 } |snapshots还可以与history做 join。例如下面这条查询会显示表历史并带上写入每个快照的 application IDselect h.made_current_at, s.operation, h.snapshot_id, h.is_current_ancestor, s.summary[spark.app.id] from prod.db.table.history h join prod.db.table.snapshots s on h.snapshot_id s.snapshot_id order by made_current_at;| made_current_at | operation | snapshot_id | is_current_ancestor | summary[spark.app.id] | | -- | -- | -- | -- | -- | | 2019-02-08 03:29:51.215 | append | 57897183625154 | true | application_1520379288616_155055 | | 2019-02-09 16:24:30.13 | delete | 29641004024753 | false | application_1520379288616_151109 | | 2019-02-09 16:32:47.336 | append | 57897183625154 | true | application_1520379288616_155055 | | 2019-02-08 03:47:55.948 | overwrite | 51792995261850 | true | application_1520379288616_152431 |Entries清单条目显示表当前所有数据文件和删除文件的清单条目SELECT * FROM prod.db.table.entries;| status | snapshot_id | sequence_number | file_sequence_number | data_file | readable_metrics | | -- | -- | -- | -- | -- | -- | | 2 | 57897183625154 | 0 | 0 | {content:0,file_path:s3:/.../table/data/00047-25-833044d0-127b-415c-b874-038a4f978c29-00612.parquet,...} | {c1:{column_size:103,value_count:15,...}} |注意entries表的列对应 manifest entry 字段status用于跟踪文件的添加和删除snapshot_id文件被添加或移除时所在快照的 IDsequence_number用于跨快照的变更排序file_sequence_number指示文件何时被添加data_file包含数据文件元数据的 struct见 data file 字段readable_metrics列提供从data_file列派生的人类可读列级指标 map便于检查与调试文件级统计信息。Files文件显示表当前的数据文件以及删除文件SELECT * FROM prod.db.table.files;返回content、file_path、file_format、spec_id、record_count、file_size_in_bytes、column_sizes、value_counts、null_value_counts、nan_value_counts、lower_bounds、upper_bounds、key_metadata、split_offsets、equality_ids、sort_order_id、readable_metrics等列。示例中既包含普通数据文件content 0也包含位置删除文件content 1*-deletes.parquet与等值删除文件content 2。说明content指数据文件存储的内容类型0 - Data数据、1 - Position Deletes位置删除、2 - Equality Deletes等值删除。只查看数据文件或删除文件分别查询prod.db.table.data_files和prod.db.table.delete_files。要跨所有被追踪的快照查看全部文件、数据文件与删除文件分别查询prod.db.table.all_files、prod.db.table.all_data_files和prod.db.table.all_delete_files。Manifests清单文件显示表当前的文件清单SELECT * FROM prod.db.table.manifests;| content | path | length | partition_spec_id | added_snapshot_id | added_data_files_count | existing_data_files_count | deleted_data_files_count | added_delete_files_count | existing_delete_files_count | deleted_delete_files_count | partition_summaries | | -- | -- | -- | -- | -- | -- | -- | -- | -- | -- | -- | -- | | 0 | s3://.../table/metadata/45b5290b-ee61-4788-b324-b1e2735c0e10-m0.avro | 4479 | 0 | 6668963634911763636 | 8 | 0 | 0 | 0 | 0 | 0 | [[false,null,2019-05-13,2019-05-15]] |注意manifests表的partition_summaries列内字段对应 manifest list 中的field_summarystruct顺序为contains_null、contains_nan、lower_bound、upper_bound。contains_nan可能返回 null表示该信息无法从文件元数据中获得。这通常发生在读取 V1 表时V1 不填充contains_nan。Partitions分区显示表当前的各个分区SELECT * FROM prod.db.table.partitions;partitionspec_idrecord_countfile_counttotal_data_file_size_in_bytesposition_delete_record_countposition_delete_file_countequality_delete_record_countequality_delete_file_countlast_updated_at(μs)last_updated_snapshot_id{20211001, 11}011100210016330860341920009205185327307503337{20211002, 11}04350011001633172537358000867027598972211003{20211001, 10}074700000016330825987160003280122546965981531{20211002, 10}032400001116331691594890006941468797545315876注意对未分区表partitions表不包含partition和spec_id字段。partitions元数据表显示当前快照中带有数据文件或删除文件的分区。但删除文件并未应用因此在某些情况下即使某个分区的所有数据行都被删除文件标记为已删除该分区仍会被显示。Positional Delete Files位置删除文件显示表当前快照中的所有位置删除文件SELECT * from prod.db.table.position_deletes;| file_path | pos | row | partition | spec_id | delete_file_path | | -- | -- | -- | -- | -- | -- | | s3:/.../table/data/00042-3-a9aa8b24-20bc-4d56-93b0-6b7675782bb5-00001.parquet | 1 | 0 | {20211001, 11} | 0 | s3:/.../table/data/00191-1933-25e9f2f3-d863-4a69-a5e1-f9aeeebe60bb-00001-deletes.parquet |All 系列元数据表All 系列元数据表是当前快照各元数据表的并集返回跨所有快照的元数据。危险由于元数据文件可能属于多个表快照all 系列元数据表可能对每个数据文件或清单文件产生多行结果。All Data Files显示表的所有数据文件及每个文件的元数据SELECT * FROM prod.db.table.all_data_files;| content | file_path | file_format | spec_id | partition | record_count | file_size_in_bytes | column_sizes | value_counts | null_value_counts | nan_value_counts | lower_bounds | upper_bounds | key_metadata | split_offsets | equality_ids | sort_order_id | readable_metrics | | -- | -- | -- | -- | -- | -- | -- | -- | -- | -- | -- | -- | -- | -- | -- | -- | -- | -- | | 0 | s3://.../dt20210102/00000-0-756e2512-49ae-45bb-aae3-c0ca475e7879-00001.parquet | PARQUET | 0 | {20210102} | 14 | 2444 | {1 - 94, 2 - 17} | {1 - 14, 2 - 14} | {1 - 0, 2 - 0} | {} | {1 - 1, 2 - 20210102} | {1 - 2, 2 - 20210102} | null | [4] | null | 0 | {id:{column_size:94,...}} | | 0 | s3://.../dt20210103/00000-0-26222098-032f-472b-8ea5-651a55b21210-00001.parquet | PARQUET | 0 | {20210103} | 14 | 2444 | {1 - 94, 2 - 17} | {1 - 14, 2 - 14} | {1 - 0, 2 - 0} | {} | {1 - 1, 2 - 20210103} | {1 - 3, 2 - 20210103} | null | [4] | null | 0 | {id:{column_size:94,...}} | | 0 | s3://.../dt20210104/00000-0-a3bb1927-88eb-4f1c-bc6e-19076b0d952e-00001.parquet | PARQUET | 0 | {20210104} | 14 | 2444 | {1 - 94, 2 - 17} | {1 - 14, 2 - 14} | {1 - 0, 2 - 0} | {} | {1 - 1, 2 - 20210104} | {1 - 3, 2 - 20210104} | null | [4] | null | 0 | {id:{column_size:94,...}} |All Delete Files显示所有快照中表的删除文件及每个文件的元数据SELECT * FROM prod.db.table.all_delete_files;All Entries显示所有快照中表的数据文件和删除文件的清单条目SELECT * FROM prod.db.table.all_entries;All Manifests显示表的所有清单文件SELECT * FROM prod.db.table.all_manifests;该表在manifests的基础上额外提供reference_snapshot_id列指示清单所属的快照| content | path | length | partition_spec_id | added_snapshot_id | added_data_files_count | existing_data_files_count | deleted_data_files_count | added_delete_files_count | existing_delete_files_count | deleted_delete_files_count | partition_summaries | reference_snapshot_id | | -- | -- | -- | -- | -- | -- | -- | -- | -- | -- | -- | -- | -- | | 0 | s3://.../metadata/a85f78c5-3222-4b37-b7e4-faf944425d48-m0.avro | 6376 | 0 | 6272782676904868561 | 2 | 0 | 0 | 0 | 0 | 0 | [{false, false, 20210101, 20210101}] | 57897183625154 |其partition_summaries字段含义与manifests表相同对应 manifest list 中的field_summarystruct顺序为contains_null、contains_nan、lower_bound、upper_boundcontains_nan在 V1 表中同样可能为 null。References快照引用显示表已知的快照引用SELECT * FROM prod.db.table.refs;| name | type | snapshot_id | max_reference_age_in_ms | min_snapshots_to_keep | max_snapshot_age_in_ms | | -- | -- | -- | -- | -- | -- | | main | BRANCH | 4686954189838128572 | 10 | 20 | 30 | | testTag | TAG | 4686954189838128572 | 10 | null | null |用 DataFrame 检查元数据表元数据表也可以使用 DataFrameReader API 加载// 命名元存储表 spark.read.format(iceberg).load(db.table.files) // Hadoop 路径表 spark.read.format(iceberg).load(hdfs://nn:8020/path/to/table#files)从源码看这些元数据表在 core 包中各自独立实现例如 FilesTable.java 继承自BaseFilesTable其manifests()通过snapshot().allManifests(table().io())枚举当前快照的全部清单同目录下还有 SnapshotsTable.java、HistoryTable.java、ManifestsTable.java、PartitionsTable.java、RefsTable.java、PositionDeletesTable.java、MetadataLogEntriesTable.java、AllDataFilesTable.java、AllManifestsTable.java、AllEntriesTable.java 等均以把表的元数据变成可查询的行为设计目标。元数据表的时态旅行元数据表同样支持时态旅行-- 获取 2021-09-20 08:00:00 时刻表的文件清单 SELECT * FROM prod.db.table.manifests TIMESTAMP AS OF 2021-09-20 08:00:00; -- 获取快照 ID 为 10963874102873L 的分区 SELECT * FROM prod.db.table.partitions VERSION AS OF 10963874102873;元数据表也可以结合 DataFrameReader API 进行时态旅行检查// 以 DataFrame 形式加载快照 ID 为 10963874102873 时的文件元数据 spark.read.format(iceberg).option(snapshot-id, 10963874102873L).load(db.table.files)实战要点小结查询入口SQL 使用catalog.namespace.table三段式标识符DataFrame 使用spark.table(...)或spark.read.format(iceberg).load(...)加载形式按匹配 catalog 优先于命名空间解析的优先级判定。变换函数system.bucket/years/months/days/hours/truncate只通过 Iceberg catalog 的system命名空间暴露返回 Iceberg 变换值如自 1970 起算的年/月/小时数适合与分区变换对齐编写过滤器。时态旅行SQL 用TIMESTAMP/VERSION AS OF与FOR SYSTEM_TIME/VERSION AS OFDataFrame 用snapshot-id、as-of-timestamp、branch、tag四个 option注意分支查询用表 schema、标签查询用快照 schema 的差异以及分支/标签标识符需用反引号转义。增量读取仅 DataFrame API 支持通过start-snapshot-id不含与end-snapshot-id含限定范围只能读取 append 操作追加的数据。表诊断history、metadata_log_entries、snapshots、entries、files、manifests、partitions、position_deletes、refs及all_*系列元数据表配合时态旅行子句可精确还原任意时间点的表状态并借助readable_metrics列调试文件级统计信息。如需了解更底层的读取行为可继续阅读 spark-configuration.md 中关于 read options、scan 规划与向量化读取的说明以及仓库中对应的源码与测试读取选项常量定义在 SparkReadOptions.java扫描构建逻辑见 SparkScanBuilder.java元数据表实现位于 core/src/main/java/org/apache/iceberg相关行为断言可在 TestDataSourceOptions.java 与 TestSparkFunctions.java 中查阅。赞分享数据湖大数据数据存储【免费下载链接】icebergApache Iceberg项目地址https://gitcode.com/gh_mirrors/icebe/iceberg点击查看免费下载相关推荐Apache Iceberg Flink 查询指南SQL 与 DataStream 的批量读取、流式读取与元数据表查询Apache Iceberg Flink 查询指南SQL 与 DataStream 的批量读取、流式读取与元数据表查询 Apache Iceberg 与 Ap数据湖大数据数据存储Cloudflare R2 SQL 实战指南用 Wrangler 与 HTTP API 查询 Apache Iceberg 表Cloudflare R2 SQL 实战指南用 Wrangler 与 HTTP API 查询 Apache Iceberg 表 R2 SQL 是 Cloudf人工智能语音音频NLP媒体生成Apache Iceberg时间旅行功能查询历史数据的完整教程Apache Iceberg时间旅行功能查询历史数据的完整教程 Apache Iceberg时间旅行功能让您能够轻松查询和分析历史数据版本就像穿越到过去查看数据湖湖仓一体大数据上一篇三步让老旧Mac重获新生OpenCore Legacy Patcher终极指南下一篇终极Mac鼠标优化指南让普通鼠标超越苹果触控板的5个专业技巧创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表