Apache Spark

查询

师成师成· 更新于 2026-09-28· 阅读 72 分钟· 0 次阅读

登录后可跨设备保存划线和私人笔记登录

Spark 查询

要在 Spark 中使用 Iceberg,首先需要配置 Spark catalog。Iceberg 使用 Apache Spark 的 DataSourceV2 API 实现数据源与 catalog。

使用 SQL 查询

在 Spark 中,表使用包含 catalog 名称的标识符。

SELECT * FROM prod.db.table; -- catalog: prod, namespace: db, table: table

元数据表(如 history 和 snapshots)可以使用 Iceberg 表名作为命名空间。

例如,要读取 prod.db.table 的 files 元数据表:

SELECT * FROM prod.db.table.files;
contentfile_pathfile_formatspec_idpartitionrecord_countfile_size_in_bytescolumn_sizesvalue_countsnull_value_countsnan_value_countslower_boundsupper_boundskey_metadatasplit_offsetsequality_idssort_order_id
0s3:/.../table/data/00000-3-8d6d60e8-d427-4809-bcf0-f5d45a4aad96.parquetPARQUET0{1999-01-01, 01}1597[1 -> 90, 2 -> 62][1 -> 1, 2 -> 1][1 -> 0, 2 -> 0][][1 -> , 2 -> c][1 -> , 2 -> c]null[4]nullnull
0s3:/.../table/data/00001-4-8d6d60e8-d427-4809-bcf0-f5d45a4aad96.parquetPARQUET0{1999-01-01, 02}1597[1 -> 90, 2 -> 62][1 -> 1, 2 -> 1][1 -> 0, 2 -> 0][][1 -> , 2 -> b][1 -> , 2 -> b]null[4]nullnull
0s3:/.../table/data/00002-5-8d6d60e8-d427-4809-bcf0-f5d45a4aad96.parquetPARQUET0{1999-01-01, 03}1597[1 -> 90, 2 -> 62][1 -> 1, 2 -> 1][1 -> 0, 2 -> 0][][1 -> , 2 -> a][1 -> , 2 -> a]null[4]nullnull

Spark SQL 函数

Iceberg 会为每个 Iceberg catalog 添加 SQL 函数,以便在查询中查看转换结果,并编写与 Iceberg 分区转换相匹配的过滤条件。这些函数只能通过 Iceberg catalog 使用;它们不会注册到 Spark 的内置 catalog 中。

注意

4.2.0 之前的 Spark 不支持会话 catalog 中的 V2Function。即使使用 org.apache.iceberg.spark.SparkSessionCatalog 配置了 spark_catalog,SELECT 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 输入,原值不变地返回;对于时间戳,则丢弃时间部分)。

对于数值输入,system.truncate(width, col) 会向下取整到最接近的 width 的倍数。对于 string 和 binary 输入,它会保留前 width 个字符或字节。

当你想要查看 Iceberg 如何转换值,或者在编写与分区转换相匹配的查询和行级操作的过滤条件时,这些函数尤其有用。

使用 SQL 进行时间旅行查询

Spark 在 SQL 查询中支持使用 TIMESTAMP AS OF 或 VERSION AS OF 子句进行时间旅行。VERSION AS OF 子句可以包含一个长整型快照 ID,或一个字符串分支名或标签名。

Info

注意:如果分支或标签的名称与某个快照 ID 相同,那么用于时间旅行的快照将是具有该给定快照 ID 的快照。例如,考虑存在一个名为 '1' 的标签,它引用了 ID 为 2 的快照。如果版本旅行子句是 VERSION AS OF '1',时间旅行将指向 ID 为 1 的快照。如果这不是期望的行为,请使用明确定义的前缀重命名该标签或分支,例如 'snapshot-1'。

-- time travel to October 26, 1986 at 01:21:00
SELECT * FROM prod.db.table TIMESTAMP AS OF '1986-10-26 01:21:00';

-- time travel to snapshot with id 10963874102873L
SELECT * FROM prod.db.table VERSION AS OF 10963874102873;

-- time travel to the head snapshot of audit-branch
SELECT * FROM prod.db.table VERSION AS OF 'audit-branch';

-- time travel to the snapshot referenced by the tag 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 时间戳(以秒为单位)的形式提供:

-- timestamp in seconds
SELECT * FROM prod.db.table TIMESTAMP AS OF 499162860;
SELECT * FROM prod.db.table FOR SYSTEM_TIME AS OF 499162860;

分支或标签也可以采用与元数据表类似的语法指定,即 branch_<branchname> 或 tag_<tagname>:

SELECT * FROM prod.db.table.`branch_audit-branch`;
SELECT * FROM prod.db.table.`tag_historical-snapshot`;

(含 "-" 的标识符无效,因此必须使用反引号进行转义。)

请注意,带分支或标签的标识符不能与 VERSION AS OF 结合使用。

时间旅行查询中的 Schema 选择

上一节中提到的各种时间旅行查询,既可以使用快照的 schema,也可以使用表的 schema:

-- time travel to October 26, 1986 at 01:21:00 -> uses the snapshot's schema
SELECT * FROM prod.db.table TIMESTAMP AS OF '1986-10-26 01:21:00';

-- time travel to snapshot with id 10963874102873L -> uses the snapshot's schema
SELECT * FROM prod.db.table VERSION AS OF 10963874102873;

-- time travel to the head of audit-branch -> uses the table's schema
SELECT * FROM prod.db.table VERSION AS OF 'audit-branch';
SELECT * FROM prod.db.table.`branch_audit-branch`;

-- time travel to the snapshot referenced by the tag historical-snapshot -> uses the snapshot's schema
SELECT * FROM prod.db.table VERSION AS OF 'historical-snapshot';
SELECT * FROM prod.db.table.`tag_historical-snapshot`;

例如,考虑一个随时间推移不断演进其 schema 的表,看看每种时间旅行查询类型是如何选择其 schema 的:

-- snapshot S1: initial schema (id, status)
CREATE TABLE prod.db.orders (
  id BIGINT,
  status STRING
) USING iceberg;

INSERT INTO prod.db.orders VALUES (1, 'NEW'), (2, 'PAID');

-- record snapshot S1's snapshot_id and committed_at timestamp
-- e.g. snapshot_id = 101, committed_at = '2025-01-01 10:00:00'

-- snapshot S2: add a new column "total" and write new data
ALTER TABLE prod.db.orders ADD COLUMN total DOUBLE;

INSERT INTO prod.db.orders VALUES (3, 'PAID', 100.0);

-- now S2 is the current snapshot with schema (id, status, total)

时间旅行查询在选定特定快照或时间戳时,会使用该快照的 schema:

-- uses the snapshot schema of S1: columns (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 模式中并不存在,因此不可见——尽管当前表模式中包含 total。

现在创建一个分支和一个标签,两者都引用 S1:

-- branch "audit_branch" points to snapshot S1
ALTER TABLE prod.db.orders CREATE BRANCH audit_branch AS OF VERSION 101;

-- tag "first_load" also points to snapshot S1
ALTER TABLE prod.db.orders CREATE TAG first_load AS OF VERSION 101;

当你查询分支时,Spark 会使用表的当前 schema:

-- uses the table schema: columns (id, status, total)
SELECT * FROM prod.db.orders VERSION AS OF 'audit_branch';

-- equivalent identifier form
SELECT * FROM prod.db.orders.`branch_audit_branch`;

在这些查询中,结果包含 (id, status, total) 三列。对于来自 S1 的行,total 被返回为 NULL,因为这些行被写入时该列尚不存在。

查询标签(tag)时,Spark 会使用该标签所引用的快照的 schema:

-- uses the snapshot schema of S1: columns (id, status)
SELECT * FROM prod.db.orders VERSION AS OF 'first_load';

-- equivalent identifier form
SELECT * FROM prod.db.orders.`tag_first_load`;

这些查询只返回 id 和 status,因为标签绑定到特定的快照,并使用该快照的 schema,即使表的当前 schema 已经发生演进。

使用 DataFrame 进行查询

要将表加载为 DataFrame,请使用 table:

val df = spark.table("prod.db.table")

使用 DataFrameReader 访问目录

路径和表名可以通过 Spark 的 DataFrameReader 接口加载。表的加载方式取决于标识符的指定方式。使用 spark.read.format("iceberg").load(table) 或 spark.table(table) 时,table 变量可以采用以下几种形式:

  • file:///path/to/table:加载指定路径下的 HadoopTable
  • tablename:加载 currentCatalog.currentNamespace.tablename
  • catalog.tablename:从指定的 catalog 加载 tablename
  • namespace.tablename:从当前 catalog 加载 namespace.tablename
  • catalog.namespace.tablename:从指定的 catalog 加载 namespace.tablename
  • namespace1.namespace2.tablename:从当前 catalog 加载 namespace1.namespace2.tablename

以上列表按优先级顺序排列。例如:匹配到的 catalog 的优先级高于任何命名空间解析。

使用 DataFrame 进行时间旅行查询

若要在 DataFrame API 中选择特定的表快照,或选择某个时间点的快照,Iceberg 支持四种 Spark 读取选项:

  • snapshot-id:选择特定的表快照
  • as-of-timestamp:选择指定时间戳(毫秒)处的当前快照
  • branch:选择指定分支的头部快照。注意,目前 branch 不能与 as-of-timestamp 组合使用。
  • tag:选择与指定标签关联的快照。tag 不能与 as-of-timestamp 组合使用。
// time travel to October 26, 1986 at 01:21:00
spark.read
    .option("as-of-timestamp", "499162860000")
    .format("iceberg")
    .load("path/to/table")
// time travel to snapshot with ID 10963874102873L
spark.read
    .option("snapshot-id", 10963874102873L)
    .format("iceberg")
    .load("path/to/table")
// time travel to tag historical-snapshot
spark.read
    .option(SparkReadOptions.TAG, "historical-snapshot")
    .format("iceberg")
    .load("path/to/table")
// time travel to the head snapshot of audit-branch
spark.read
    .option(SparkReadOptions.BRANCH, "audit-branch")
    .format("iceberg")
    .load("path/to/table")

增量读取

要增量读取追加的数据,请使用:

  • start-snapshot-id 增量扫描使用的起始快照 ID(不包含该快照)。
  • end-snapshot-id 增量扫描使用的结束快照 ID(包含该快照)。此参数为可选;若省略,则默认为当前快照。
// get the data added after start-snapshot-id (10963874102873L) until 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 格式版本。Spark 的 SQL 语法不支持增量读取。

检查表

要检查表的历史、快照和其他元数据,Iceberg 支持元数据表。

元数据表通过在原表名后追加元数据表名来标识。例如,db.table 的历史通过 db.table.history 读取。

历史

要显示表的历史:

SELECT * FROM prod.db.table.history;
made_current_atsnapshot_idparent_idis_current_ancestor
2019-02-08 03:29:51.2155781947118336215154NULLtrue
2019-02-08 03:47:55.94851792995261850568305781947118336215154true
2019-02-09 16:24:30.132964100402475335445179299526185056830false
2019-02-09 16:32:47.33629998756080624373305179299526185056830true
2019-02-09 19:42:03.91989245587860605834792999875608062437330true
2019-02-09 19:49:16.34365367338231819750458924558786060583479true

信息

这里显示了一个被回滚的提交。 示例中有两个快照具有相同的父级,其中一个不是当前表状态的祖先。

元数据日志条目

要显示表的元数据日志条目:

SELECT * from prod.db.table.metadata_log_entries;
timestampfilelatest_snapshot_idlatest_schema_idlatest_sequence_number
2022-07-28 10:43:52.93s3://.../table/metadata/00000-9441e604-b3c2-498a-a45a-6320e8ab9006.metadata.jsonnullnullnull
2022-07-28 10:43:57.487s3://.../table/metadata/00001-f30823df-b745-4a0a-b293-7532e0c99986.metadata.json17026083367764530001
2022-07-28 10:43:58.25s3://.../table/metadata/00002-2cc2837a-02dc-4687-acc1-b4d86ea486f4.metadata.json95890649397670977402

快照

要显示表的有效快照:

SELECT * FROM prod.db.table.snapshots;
committed_atsnapshot_idparent_idoperationmanifest_listsummary
2019-02-08 03:29:51.21557897183625154nullappends3://.../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 }

你还可以将快照与表历史进行关联。例如,下面这个查询会显示表历史,以及写入每个快照的应用 ID:

select
    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_atoperationsnapshot_idis_current_ancestorsummary[spark.app.id]
2019-02-08 03:29:51.215append57897183625154trueapplication_1520379288616_155055
2019-02-09 16:24:30.13delete29641004024753falseapplication_1520379288616_151109
2019-02-09 16:32:47.336append57897183625154trueapplication_1520379288616_155055
2019-02-08 03:47:55.948overwrite51792995261850trueapplication_1520379288616_152431

条目(Entries)

显示表中数据文件和删除文件的当前所有清单条目。

SELECT * FROM prod.db.table.entries;
statussnapshot_idsequence_numberfile_sequence_numberdata_filereadable_metrics
25789718362515400{"content":0,"file_path":"s3:/.../table/data/00047-25-833044d0-127b-415c-b874-038a4f978c29-00612.parquet","file_format":"PARQUET","spec_id":0,"record_count":15,"file_size_in_bytes":473,"column_sizes":{1:103},"value_counts":{1:15},"null_value_counts":{1:0},"nan_value_counts":{},"lower_bounds":{1:},"upper_bounds":{1:},"key_metadata":null,"split_offsets":[4],"equality_ids":null,"sort_order_id":0}{"c1":{"column_size":103,"value_count":15,"null_value_count":0,"nan_value_count":null,"lower_bound":1,"upper_bound":3}}

注意:

  1. entries 表中的各列对应于清单条目字段:

    • status:用于跟踪新增和删除操作
    • snapshot_id:文件被添加或删除所在的快照 ID
    • sequence_number:用于跨快照对变更进行排序
    • file_sequence_number:表示文件被添加的时间
    • data_file:一个包含数据文件元数据的结构体,参见数据文件字段
  2. readable_metrics 列提供了从 data_file 列派生的、便于人类阅读的列级扩展指标映射,便于检查和调试文件级统计信息。

文件

要显示表的当前文件:

SELECT * FROM prod.db.table.files;
内容file_pathfile_formatspec_idrecord_countfile_size_in_bytescolumn_sizesvalue_countsnull_value_countsnan_value_countslower_boundsupper_boundskey_metadatasplit_offsetsequality_idssort_order_idreadable_metrics
0s3:/.../table/data/00042-3-a9aa8b24-20bc-4d56-93b0-6b7675782bb5-00001.parquetPARQUET01652{1:52,2:48}{1:1,2:1}{1:0,2:0}{}{1:,2:d}{1:,2:d}NULL[4]NULL0{"data":{"column_size":48,"value_count":1,"null_value_count":0,"nan_value_count":null,"lower_bound":"d","upper_bound":"d"},"id":{"column_size":52,"value_count":1,"null_value_count":0,"nan_value_count":null,"lower_bound":1,"upper_bound":1}}
0s3:/.../table/data/00000-0-f9709213-22ca-4196-8733-5cb15d2afeb9-00001.parquetPARQUET01643{1:46,2:48}{1:1,2:1}{1:0,2:0}{}{1:,2:a}{1:,2:a}NULL[4]NULL0{"data":{"column_size":48,"value_count":1,"null_value_count":0,"nan_value_count":null,"lower_bound":"a","upper_bound":"a"},"id":{"column_size":46,"value_count":1,"null_value_count":0,"nan_value_count":null,"lower_bound":1,"upper_bound":1}}
0s3:/.../table/data/00001-1-f9709213-22ca-4196-8733-5cb15d2afeb9-00001.parquetPARQUET02644{1:49,2:51}{1:2,2:2}{1:0,2:0}{}{1:,2:b}{1:,2:c}NULL[4]NULL0{"data":{"column_size":51,"value_count":2,"null_value_count":0,"nan_value_count":null,"lower_bound":"b","upper_bound":"c"},"id":{"column_size":49,"value_count":2,"null_value_count":0,"nan_value_count":null,"lower_bound":2,"upper_bound":3}}

| 1 | s3:/.../table/data/00081-4-a9aa8b24-20bc-4d56-93b0-6b7675782bb5-00001-deletes.parquet | PARQUET | 0 | 1 | 1560 | {2147483545:46,2147483546:152} | {2147483545:1,2147483546:1} | {2147483545:0,2147483546:0} | {} | {2147483545:,2147483546:s3:/.../table/data/00000-0-f9709213-22ca-4196-8733-5cb15d2afeb9-00001.parquet} | {2147483545:,2147483546:s3:/.../table/data/00000-0-f9709213-22ca-4196-8733-5cb15d2afeb9-00001.parquet} | NULL | [4] | NULL | NULL | {"

Info

Content 表示数据文件所存储的内容类型:

  • 0 - 数据
  • 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

要显示表当前的文件清单(manifest):

SELECT * FROM prod.db.table.manifests;
contentpathlengthpartition_spec_idadded_snapshot_idadded_data_files_countexisting_data_files_countdeleted_data_files_countadded_delete_files_countexisting_delete_files_countdeleted_delete_files_countpartition_summaries
0s3://.../table/metadata/45b5290b-ee61-4788-b324-b1e2735c0e10-m0.avro447906668963634911763636800000[[false,null,2019-05-13,2019-05-15]]

注意:

  1. manifests 表中 partition_summaries 列内的各个字段对应于清单列表中的 field_summary 结构体,顺序如下:

    • contains_null
    • contains_nan
    • lower_bound
    • upper_bound
  2. contains_nan 可能返回 null,表示无法从文件的元数据中获取该信息。这种情况通常出现在读取 V1 表时,因为 V1 表不会填充 contains_nan。

分区

要显示表的当前分区:

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

注意:

  1. 对于未分区的表,partitions 元数据表中不会包含 partition 和 spec_id 字段。
  2. partitions 元数据表显示当前快照中包含数据文件或删除文件的分区。不过,删除文件并未被应用,因此在某些情况下,即使某个分区的所有数据行都已被删除文件标记为删除,该分区仍可能被显示出来。

位置删除文件

要显示表当前快照中的所有位置删除文件:

SELECT * from prod.db.table.position_deletes;
file_pathposrowpartitionspec_iddelete_file_path
s3:/.../table/data/00042-3-a9aa8b24-20bc-4d56-93b0-6b7675782bb5-00001.parquet10{20211001, 11}0s3:/.../table/data/00191-1933-25e9f2f3-d863-4a69-a5e1-f9aeeebe60bb-00001-deletes.parquet

全部元数据表

这些表是当前快照所对应各元数据表的并集,返回所有快照范围内的元数据。

危险

"all" 元数据表可能会针对同一个数据文件或 manifest 文件返回多行,因为同一个元数据文件可能属于多个表快照。

全部数据文件

显示表中的所有数据文件及其各自的元数据:

SELECT * FROM prod.db.table.all_data_files;
contentfile_pathfile_formatspec_idpartitionrecord_countfile_size_in_bytescolumn_sizesvalue_countsnull_value_countsnan_value_countslower_boundsupper_boundskey_metadatasplit_offsetsequality_idssort_order_idreadable_metrics
0s3://.../dt=20210102/00000-0-756e2512-49ae-45bb-aae3-c0ca475e7879-00001.parquetPARQUET0{20210102}142444{1 -> 94, 2 -> 17}{1 -> 14, 2 -> 14}{1 -> 0, 2 -> 0}{}{1 -> 1, 2 -> 20210102}{1 -> 2, 2 -> 20210102}null[4]null0{"id":{"column_size":94,"value_count":14,"null_value_count":0,"nan_value_count":null,"lower_bound":1,"upper_bound":2},"data":{"column_size":17,"value_count":14,"null_value_count": 0,"nan_value_count":null,"lower_bound":20210102,"upper_bound":20210102}}
0s3://.../dt=20210103/00000-0-26222098-032f-472b-8ea5-651a55b21210-00001.parquetPARQUET0{20210103}142444{1 -> 94, 2 -> 17}{1 -> 14, 2 -> 14}{1 -> 0, 2 -> 0}{}{1 -> 1, 2 -> 20210103}{1 -> 3, 2 -> 20210103}null[4]null0{"id":{"column_size":94,"value_count":14,"null_value_count":0,"nan_value_count":null,"lower_bound":1,"upper_bound":3},"data":{"column_size":17,"value_count":14,"null_value_count": 0,"nan_value_count":null,"lower_bound":20210103,"upper_bound":20210103}}
0s3://.../dt=20210104/00000-0-a3bb1927-88eb-4f1c-bc6e-19076b0d952e-00001.parquetPARQUET0{20210104}142444{1 -> 94, 2 -> 17}{1 -> 14, 2 -> 14}{1 -> 0, 2 -> 0}{}{1 -> 1, 2 -> 20210104}{1 -> 3, 2 -> 20210104}null[4]null0{"id":{"column_size":94,"value_count":14,"null_value_count":0,"nan_value_count":null,"lower_bound":1,"upper_bound":3},"data":{"column_size":17,"value_count":14,"null_value_count": 0,"nan_value_count":null,"lower_bound":20210104,"upper_bound":20210104}}

所有删除文件

要显示表中的删除文件以及所有快照中每个文件的元数据:

SELECT * FROM prod.db.table.all_delete_files;
contentfile_pathfile_formatspec_idpartitionrecord_countfile_size_in_bytescolumn_sizesvalue_countsnull_value_countsnan_value_countslower_boundsupper_boundskey_metadatasplit_offsetsequality_idssort_order_idreadable_metrics
1s3:/.../table/data/00081-4-a9aa8b24-20bc-4d56-93b0-6b7675782bb5-00001-deletes.parquetPARQUET0{20210102}11560{2147483545:46,2147483546:152}{2147483545:1,2147483546:1}{2147483545:0,2147483546:0}{}{2147483545:,2147483546:s3:/.../table/data/00000-0-f9709213-22ca-4196-8733-5cb15d2afeb9-00001.parquet}{2147483545:,2147483546:s3:/.../table/data/00000-0-f9709213-22ca-4196-8733-5cb15d2afeb9-00001.parquet}NULL[4]NULLNULL{"data":{"column_size":null,"value_count":null,"null_value_count":null,"nan_value_count":null,"lower_bound":null,"upper_bound":null},"id":{"column_size":null,"value_count":null,"null_value_count":null,"nan_value_count":null,"lower_bound":null,"upper_bound":null}}
2s3:/.../table/data/00047-25-833044d0-127b-415c-b874-038a4f978c29-00612.parquetPARQUET0{20210103}12650628613985{100:135377,101:11314}{100:126506,101:126506}{100:105434,101:11}{}{100:0,101:17}{100:404455227527,101:23}NULLNULL[1]0{"id":{"column_size":135377,"value_count":126506,"null_value_count":105434,"nan_value_count":null,"lower_bound":0,"upper_bound":404455227527},"data":{"column_size":11314,"value_count":126506,"null_value_count": 11,"nan_value_count":null,"lower_bound":17,"upper_bound":23}}

所有条目

要显示所有快照中数据文件和删除文件的清单条目:

SELECT * FROM prod.db.table.all_entries;
statussnapshot_idsequence_numberfile_sequence_numberdata_filereadable_metrics
25789718362515400{"content":0,"file_path":"s3:/.../table/data/00047-25-833044d0-127b-415c-b874-038a4f978c29-00612.parquet","file_format":"PARQUET","spec_id":0,"record_count":15,"file_size_in_bytes":473,"column_sizes":{1:103},"value_counts":{1:15},"null_value_counts":{1:0},"nan_value_counts":{},"lower_bounds":{1:},"upper_bounds":{1:},"key_metadata":null,"split_offsets":[4],"equality_ids":null,"sort_order_id":0}{"c1":{"column_size":103,"value_count":15,"null_value_count":0,"nan_value_count":null,"lower_bound":1,"upper_bound":3}}

所有清单文件

要显示该表的所有清单文件:

SELECT * FROM prod.db.table.all_manifests;
contentpathlengthpartition_spec_idadded_snapshot_idadded_data_files_countexisting_data_files_countdeleted_data_files_countadded_delete_files_countexisting_delete_files_countdeleted_delete_files_countpartition_summariesreference_snapshot_id
0s3://.../metadata/a85f78c5-3222-4b37-b7e4-faf944425d48-m0.avro637606272782676904868561200000[{false, false, 20210101, 20210101}]57897183625154

注意:

  1. manifests 表中 partition_summaries 列内的字段对应于 manifest list 中的 field_summary 结构体,顺序如下:

    • contains_null
    • contains_nan
    • lower_bound
    • upper_bound
  2. contains_nan 可能返回 null,这表示该信息在文件的元数据中不可用。这种情况通常出现在读取 V1 表时,因为 V1 表不会填充 contains_nan。

References

显示表已知的快照引用:

SELECT * FROM prod.db.table.refs;
nametypesnapshot_idmax_reference_age_in_msmin_snapshots_to_keepmax_snapshot_age_in_ms
mainBRANCH4686954189838128572102030
testTagTAG468695418983812857210nullnull

使用 DataFrame 检查

可以使用 DataFrameReader API 加载元数据表:

// named metastore table
spark.read.format("iceberg").load("db.table.files")
// Hadoop path table
spark.read.format("iceberg").load("hdfs://nn:8020/path/to/table#files")

使用元数据表进行时间旅行

要使用时间旅行功能检查表的元数据:

-- get the table's file manifests at timestamp Sep 20, 2021 08:00:00
SELECT * FROM prod.db.table.manifests TIMESTAMP AS OF '2021-09-20 08:00:00';

-- get the table's partitions with snapshot id 10963874102873L
SELECT * FROM prod.db.table.partitions VERSION AS OF 10963874102873;

元数据表同样可以通过 DataFrameReader API 使用时间旅行功能来查看:

// load the table's file metadata at snapshot-id 10963874102873 as DataFrame
spark.read.format("iceberg").option("snapshot-id", 10963874102873L).load("db.table.files")

评论

登录后参与评论

正在加载评论…