Apache Spark

DDL

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

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

Spark DDL

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

CREATE TABLE

Spark 3 可以通过 USING iceberg 子句在任意 Iceberg catalog 中创建表:

CREATE TABLE prod.db.sample (
    id bigint NOT NULL COMMENT 'unique id',
    data string)
USING iceberg;

Iceberg 会将 Spark 中的列类型转换为对应的 Iceberg 类型。详情请查看 建表时的类型兼容性 一节。

建表命令(包括 CTAS 和 RTAS)支持完整的 Spark 建表子句,包括:

  • PARTITIONED BY (partition-expressions) 用于配置分区
  • LOCATION '(fully-qualified-uri)' 用于设置表的位置
  • COMMENT 'table documentation' 用于设置表的描述
  • TBLPROPERTIES ('key'='value', ...) 用于设置表配置

建表命令还可以通过 USING 子句设置默认格式。此功能仅适用于 SparkCatalog,因为 Spark 对内置目录的 USING 子句处理方式不同。

不支持 CREATE TABLE ... LIKE ... 语法。

PARTITIONED BY

要创建分区表,请使用 PARTITIONED BY:

CREATE TABLE prod.db.sample (
    id bigint,
    data string,
    category string)
USING iceberg
PARTITIONED BY (category);

PARTITIONED BY 子句支持转换表达式,用于创建隐藏分区。

CREATE TABLE prod.db.sample (
    id bigint,
    data string,
    category string,
    ts timestamp)
USING iceberg
PARTITIONED BY (bucket(16, id), days(ts), category);

支持的转换如下:

  • year(ts):按年分区

  • month(ts):按月分区

  • day(ts) 或 date(ts):等同于 dateint 分区

  • hour(ts) 或 date_hour(ts):等同于 dateint 和小时分区

  • bucket(N, col):按哈希值对 N 个桶取模进行分区

  • truncate(L, col):按截断到 L 的值进行分区

    • 字符串被截断到给定长度
    • 整数和长整型被截断到对应的分桶:truncate(10, i) 会产生 0、10、20、30…… 等分区

注意:为了兼容性,旧语法 years(ts)、months(ts)、days(ts) 和 hours(ts) 同样受到支持。

相同的转换也可作为 Spark SQL 函数在 system 命名空间下使用。参见 Spark SQL 函数。

CREATE TABLE ... AS SELECT

当使用 SparkCatalog 时,Iceberg 支持将 CTAS 作为原子操作。使用 SparkSessionCatalog 时也支持 CTAS,但不具备原子性。

CREATE TABLE prod.db.sample
USING iceberg
AS SELECT ...

新创建的表不会从 SELECT 中的源表继承分区规范和表属性,你可以在 CTAS 中使用 PARTITIONED BY 和 TBLPROPERTIES 来为新表声明分区规范和表属性。

CREATE TABLE prod.db.sample
USING iceberg
PARTITIONED BY (part)
TBLPROPERTIES ('key'='value')
AS SELECT ...

REPLACE TABLE ... AS SELECT

使用 SparkCatalog 时,Iceberg 支持将 RTAS 作为原子操作。使用 SparkSessionCatalog 时也支持 RTAS,但它不是原子的。

原子表替换会以 SELECT 查询的结果创建一个新快照,同时保留表的历史记录。

REPLACE TABLE prod.db.sample
USING iceberg
AS SELECT ...
REPLACE TABLE prod.db.sample
USING iceberg
PARTITIONED BY (part)
TBLPROPERTIES ('key'='value')
AS SELECT ...
CREATE OR REPLACE TABLE prod.db.sample
USING iceberg
AS SELECT ...

如果 schema 和分区规则发生了变化,则会被替换。若要避免修改表的 schema 和分区方式,请使用 INSERT OVERWRITE 而不是 REPLACE TABLE。REPLACE TABLE 命令中的新表属性会与现有的表属性合并。现有的表属性如果发生变更则会被更新,否则保持不变。

DROP TABLE

删除表的行为在 0.14 版本中发生了变化。

在 0.14 版本之前,执行 DROP TABLE 会将表从元数据目录中移除,并同时删除表中的数据。

从 0.14 版本起,DROP TABLE 只会将表从元数据目录中移除。若要删除表中的数据,应使用 DROP TABLE PURGE。

DROP TABLE

要从元数据目录中删除表,请执行:

DROP TABLE prod.db.sample;

DROP TABLE PURGE

要从 catalog 中删除表并清除表的内容,请执行:

DROP TABLE prod.db.sample PURGE;

ALTER TABLE

Iceberg 在 Spark 3 中提供了完整的 ALTER TABLE 支持,包括:

  • 重命名表
  • 设置或移除表属性
  • 添加、删除和重命名列
  • 添加、删除和重命名嵌套字段
  • 重新排序顶层列和嵌套 struct 字段
  • 扩宽 int、float 和 decimal 字段的类型
  • 将必需列改为可选列

此外,SQL 扩展可用于增加对分区演进(partition evolution)以及设置表写入顺序的支持。

Hive Catalog 的限制

Hive Metastore(HMS)会按位置(positionally)比较列类型来验证架构变更(hive.metastore.disallow.incompatible.col.type.changes,默认为 true)。任何会改变列位置的架构演进操作在使用 Hive catalog 时都会失败。受影响的操作包括:

  • 带有 FIRST 或 AFTER 子句的 ADD COLUMN
  • 带有 FIRST 或 AFTER 子句的 ALTER COLUMN(重新排序)
  • 对非最后一列执行 DROP COLUMN

要规避此问题,可将 HMS 的架构兼容性检查禁用,设置 hive.metastore.disallow.incompatible.col.type.changes=false:

  • 远程 HMS: 在 HMS 服务器的 hive-site.xml 中设置该属性。
  • 内嵌 HMS: 启动 Spark 时传入 --conf spark.hadoop.hive.metastore.disallow.incompatible.col.type.changes=false。

权衡: 禁用此检查后,由于 Hive Metastore 中的架构不匹配,Hive 引擎可能无法再正确读取该表。支持 Iceberg 的引擎(Spark、Flink、Trino 等)仍能正常工作,因为它们从 Iceberg 元数据而非 HMS 读取架构。

ALTER TABLE ... RENAME TO

ALTER TABLE prod.db.sample RENAME TO prod.db.new_name;

ALTER TABLE ... SET TBLPROPERTIES

ALTER TABLE prod.db.sample SET TBLPROPERTIES (
    'read.split.target-size'='268435456'
);

Iceberg 使用表属性来控制表的行为。可用属性的完整列表请参见表配置。

UNSET 用于移除属性:

ALTER TABLE prod.db.sample UNSET TBLPROPERTIES ('read.split.target-size');

SET TBLPROPERTIES 也可以用来设置表注释(描述):

ALTER TABLE prod.db.sample SET TBLPROPERTIES (
    'comment' = 'A table comment.'
);

ALTER TABLE ... ADD COLUMN

要向 Iceberg 表中添加列,请在 ALTER TABLE 中使用 ADD COLUMNS 子句:

ALTER TABLE prod.db.sample
ADD COLUMNS (
    new_column string comment 'new_column docs'
);

可以一次添加多个列,列与列之间用逗号分隔。

嵌套列应使用完整列名来标识:

-- create a struct column
ALTER TABLE prod.db.sample
ADD COLUMN point struct<x: double, y: double>;

-- add a field to the struct
ALTER TABLE prod.db.sample
ADD COLUMN point.z double;
-- create a nested array column of struct
ALTER TABLE prod.db.sample
ADD COLUMN points array<struct<x: double, y: double>>;

-- add a field to the struct within an array. Using keyword 'element' to access the array's element column.
ALTER TABLE prod.db.sample
ADD COLUMN points.element.z double;
-- create a map column of struct key and struct value
ALTER TABLE prod.db.sample
ADD COLUMN points map<struct<x: int>, struct<a: int>>;

-- add a field to the value struct in a map. Using keyword 'value' to access the map's value column.
ALTER TABLE prod.db.sample
ADD COLUMN points.value.b int;

注意:不允许通过添加列的方式修改映射(map)的 key 列,只能更新映射的 value 列。

可以通过添加 FIRST 或 AFTER 子句,在任意位置添加列:

ALTER TABLE prod.db.sample
ADD COLUMN new_column bigint AFTER other_column;
ALTER TABLE prod.db.sample
ADD COLUMN nested.new_column bigint FIRST;

Hive Catalog 的限制

使用 Hive catalog 时,通过 FIRST 或 AFTER 添加列可能会因 HMS 的位置化模式校验而失败。详情及规避方法请参阅上文的警告说明。

ALTER TABLE ... RENAME COLUMN

Iceberg 允许重命名任意字段。若要重命名字段,请使用 RENAME COLUMN:

ALTER TABLE prod.db.sample RENAME COLUMN data TO payload;
ALTER TABLE prod.db.sample RENAME COLUMN location.lat TO latitude;

请注意,嵌套的重命名命令只会重命名叶子字段。上面的命令将 location.lat 重命名为 location.latitude。

ALTER TABLE ... ALTER COLUMN

ALTER COLUMN 用于拓宽类型、将字段设为可选、设置注释以及重新排列字段顺序。

如果更新是安全的,Iceberg 允许更新列类型。安全的更新包括:

  • int 转为 bigint
  • float 转为 double
  • decimal(P,S) 转为 decimal(P2,S),其中 P2 > P(标度 S 不能改变)
ALTER TABLE prod.db.sample ALTER COLUMN measurement TYPE double;

要向 struct 中添加列或删除列,请对嵌套列名使用 ADD COLUMN 或 DROP COLUMN。

列注释也可以通过 ALTER COLUMN 更新:

ALTER TABLE prod.db.sample ALTER COLUMN measurement TYPE double COMMENT 'unit is bytes per second';
ALTER TABLE prod.db.sample ALTER COLUMN measurement COMMENT 'unit is kilobytes per second';

Iceberg 允许使用 FIRST 和 AFTER 子句重新排列顶层列或 struct 中的列:

ALTER TABLE prod.db.sample ALTER COLUMN col FIRST;
ALTER TABLE prod.db.sample ALTER COLUMN nested.col AFTER other_col;

Hive Catalog 限制

使用 Hive catalog 时,重新排列列的顺序可能会失败,因为 HMS 会进行位置模式校验。详情及变通方法请参阅上文的「Hive Catalog 限制」说明。

对于非空列,可以使用 DROP NOT NULL 来修改其可空性:

ALTER TABLE prod.db.sample ALTER COLUMN id DROP NOT NULL;

信息

无法使用 SET NOT NULL 将可空列改为非空列,因为 Iceberg 无法确定现有数据中是否包含空值。

信息

ALTER COLUMN 不用于更新 struct 类型。请使用 ADD COLUMN 和 DROP COLUMN 来添加或删除结构体字段。

ALTER TABLE ... DROP COLUMN

要删除列,请使用 ALTER TABLE ... DROP COLUMN:

ALTER TABLE prod.db.sample DROP COLUMN id;
ALTER TABLE prod.db.sample DROP COLUMN point.z;

Hive Catalog 限制

使用 Hive Catalog 时,删除非最后一列可能会因为 HMS 的位置式模式校验而失败。详情及解决方法请参见上文「Hive Catalog 限制」警告。

ALTER TABLE SQL 扩展

以下命令在 Spark 3 中使用 Iceberg SQL 扩展时可用。

ALTER TABLE ... ADD PARTITION FIELD

Iceberg 支持使用 ADD PARTITION FIELD 向分区规范中添加新的分区字段:

ALTER TABLE prod.db.sample ADD PARTITION FIELD catalog; -- identity transform

分区转换 同样受到支持:

ALTER TABLE prod.db.sample ADD PARTITION FIELD bucket(16, id);
ALTER TABLE prod.db.sample ADD PARTITION FIELD truncate(4, data);
ALTER TABLE prod.db.sample ADD PARTITION FIELD year(ts);
-- use optional AS keyword to specify a custom name for the partition field
ALTER TABLE prod.db.sample ADD PARTITION FIELD bucket(16, id) AS shard;

添加分区字段是元数据操作,不会改变任何现有表数据。新写入的数据将使用新的分区方式,而现有数据仍保持旧的分区布局。旧数据文件在元数据表中的新分区字段将为 null 值。

当表的分区方式发生变化时,动态分区覆盖的行为也会随之改变,因为动态覆盖会隐式地替换分区。若要显式覆盖,请使用新的 DataFrameWriterV2 API。

注意

若要从按天分区迁移到使用转换的按小时分区,并不必删除按天的分区字段。保留该字段可确保现有的元数据表查询继续正常工作。

危险

当分区方式改变时,动态分区覆盖的行为将会变化。 例如,如果原本按天分区并改为按小时分区,覆盖操作将覆盖小时分区,而不再覆盖天分区。

ALTER TABLE ... DROP PARTITION FIELD

可以使用 DROP PARTITION FIELD 删除分区字段:

ALTER TABLE prod.db.sample DROP PARTITION FIELD catalog;
ALTER TABLE prod.db.sample DROP PARTITION FIELD bucket(16, id);
ALTER TABLE prod.db.sample DROP PARTITION FIELD truncate(4, data);
ALTER TABLE prod.db.sample DROP PARTITION FIELD year(ts);
ALTER TABLE prod.db.sample DROP PARTITION FIELD shard;

请注意,尽管分区被移除,该列仍将保留在表结构中。

删除分区字段是一项元数据操作,不会更改任何现有表数据。新数据将按照新的分区方式写入,但现有数据将保留在旧的分区布局中。

危险

当分区方式发生变化时,动态分区覆写的行为也会随之改变。 例如,如果按天分区改为按小时分区,覆写操作将会覆写小时级分区,但不再覆写天级分区。

危险

删除分区字段时请务必小心,因为它会改变 files 等元数据表的结构,可能导致元数据查询失败或产生不同的结果。

ALTER TABLE ... REPLACE PARTITION FIELD

可以使用 REPLACE PARTITION FIELD 在单次元数据更新中将一个分区字段替换为新的分区字段:

ALTER TABLE prod.db.sample REPLACE PARTITION FIELD ts_day WITH day(ts);
-- use optional AS keyword to specify a custom name for the new partition field
ALTER TABLE prod.db.sample REPLACE PARTITION FIELD ts_day WITH day(ts) AS day_of_ts;

ALTER TABLE ... WRITE ORDERED BY

Iceberg 表可以配置排序顺序,某些引擎会据此对写入表中的数据进行自动排序。例如,Spark 中的 MERGE INTO 就会使用表的排序顺序。

要设置表的写入顺序,请使用 WRITE ORDERED BY:

ALTER TABLE prod.db.sample WRITE ORDERED BY category, id
-- use optional ASC/DESC keyword to specify sort order of each field (default ASC)
ALTER TABLE prod.db.sample WRITE ORDERED BY category ASC, id DESC
-- use optional NULLS FIRST/NULLS LAST keyword to specify null order of each field (default FIRST)
ALTER TABLE prod.db.sample WRITE ORDERED BY category ASC NULLS LAST, id DESC NULLS FIRST

信息

表的写入顺序并不保证查询时的数据顺序,它只影响数据写入表中的方式。

WRITE ORDERED BY 用于设置全局排序,使各行在所有任务之间都保持有序,类似于在 INSERT 命令中使用 ORDER BY:

INSERT INTO prod.db.sample
SELECT id, data, category, ts FROM another_table
ORDER BY ts, category

若只想在每个任务内部排序,而不在任务之间排序,请使用 LOCALLY ORDERED BY:

ALTER TABLE prod.db.sample WRITE LOCALLY ORDERED BY category, id

要取消表的排序规则,请使用 UNORDERED:

ALTER TABLE prod.db.sample WRITE UNORDERED

ALTER TABLE ... WRITE DISTRIBUTED BY PARTITION

WRITE DISTRIBUTED BY PARTITION 将要求每个分区由一个写入者处理,默认实现为哈希分布(hash distribution)。

ALTER TABLE prod.db.sample WRITE DISTRIBUTED BY PARTITION

DISTRIBUTED BY PARTITION 与 LOCALLY ORDERED BY 可以结合使用,以便按分区进行分布,并在每个任务内对行进行本地排序。

ALTER TABLE prod.db.sample WRITE DISTRIBUTED BY PARTITION LOCALLY ORDERED BY category, id

ALTER TABLE ... SET IDENTIFIER FIELDS

Iceberg 支持通过 SET IDENTIFIER FIELDS 为一个 spec 设置标识字段:如果表包含标识字段,Spark 表就可以支持 Flink SQL 的 upsert 操作。

ALTER TABLE prod.db.sample SET IDENTIFIER FIELDS id
-- single column
ALTER TABLE prod.db.sample SET IDENTIFIER FIELDS id, data
-- multiple columns

标识字段在创建或添加时必须是 NOT NULL 列。后续的 ALTER 语句会覆盖先前的设置。

ALTER TABLE ... DROP IDENTIFIER FIELDS

可以使用 DROP IDENTIFIER FIELDS 删除标识字段:

ALTER TABLE prod.db.sample DROP IDENTIFIER FIELDS id
-- single column
ALTER TABLE prod.db.sample DROP IDENTIFIER FIELDS id, data
-- multiple columns

请注意,虽然标识符被移除,但该列仍然存在于表结构中。

分支与标签 DDL

ALTER TABLE ... CREATE BRANCH

可以通过 CREATE BRANCH 语句创建分支,支持以下选项:

  • 使用 IF NOT EXISTS,当分支已存在时不会失败
  • 使用 CREATE OR REPLACE,当分支已存在时对其进行更新
  • 在特定快照处创建分支
  • 创建具有指定保留期限的分支
-- CREATE audit-branch at current snapshot with default retention.
ALTER TABLE prod.db.sample CREATE BRANCH `audit-branch`

-- CREATE audit-branch at current snapshot with default retention if it doesn't exist.
ALTER TABLE prod.db.sample CREATE BRANCH IF NOT EXISTS `audit-branch`

-- CREATE audit-branch at current snapshot with default retention or REPLACE it if it already exists.
ALTER TABLE prod.db.sample CREATE OR REPLACE BRANCH `audit-branch`

-- CREATE audit-branch at snapshot 1234 with default retention.
ALTER TABLE prod.db.sample CREATE BRANCH `audit-branch`
AS OF VERSION 1234

-- CREATE audit-branch at snapshot 1234, retain audit-branch for 30 days, and retain the latest 30 days. The latest 3 snapshot snapshots, and 2 days worth of snapshots.
ALTER TABLE prod.db.sample CREATE BRANCH `audit-branch`
AS OF VERSION 1234 RETAIN 30 DAYS
WITH SNAPSHOT RETENTION 3 SNAPSHOTS 2 DAYS

ALTER TABLE ... CREATE TAG

可以通过 CREATE TAG 语句创建标签,并支持以下选项:

  • 使用 IF NOT EXISTS,若标签已存在则不报错
  • 使用 CREATE OR REPLACE,若标签已存在则更新该标签
  • 在指定的快照处创建标签
  • 创建具有指定保留期限的标签
-- CREATE historical-tag at current snapshot with default retention.
ALTER TABLE prod.db.sample CREATE TAG `historical-tag`

-- CREATE historical-tag at current snapshot with default retention if it doesn't exist.
ALTER TABLE prod.db.sample CREATE TAG IF NOT EXISTS `historical-tag`

-- CREATE historical-tag at current snapshot with default retention or REPLACE it if it already exists.
ALTER TABLE prod.db.sample CREATE OR REPLACE TAG `historical-tag`

-- CREATE historical-tag at snapshot 1234 with default retention.
ALTER TABLE prod.db.sample CREATE TAG `historical-tag` AS OF VERSION 1234

-- CREATE historical-tag at snapshot 1234 and retain it for 1 year.
ALTER TABLE prod.db.sample CREATE TAG `historical-tag`
AS OF VERSION 1234 RETAIN 365 DAYS

ALTER TABLE ... REPLACE BRANCH

分支所引用的快照可以通过 REPLACE BRANCH SQL 进行更新。也可以在该语句中更新保留策略。

-- REPLACE audit-branch to reference snapshot 4567 and update the retention to 60 days.
ALTER TABLE prod.db.sample REPLACE BRANCH `audit-branch`
AS OF VERSION 4567 RETAIN 60 DAYS

ALTER TABLE ... REPLACE TAG

标签所引用的快照可以通过 REPLACE TAG SQL 进行更新。该语句也可以同时更新保留期(Retention)。

-- REPLACE historical-tag to reference snapshot 4567 and update the retention to 60 days.
ALTER TABLE prod.db.sample REPLACE TAG `historical-tag`
AS OF VERSION 4567 RETAIN 60 DAYS

ALTER TABLE ... DROP BRANCH

可以通过 DROP BRANCH SQL 删除分支。

ALTER TABLE prod.db.sample DROP BRANCH `audit-branch`

ALTER TABLE ... DROP TAG

标签可以通过 DROP TAG SQL 语句删除。

ALTER TABLE prod.db.sample DROP TAG `historical-tag`

Spark 中的 Iceberg 视图

Iceberg 视图是 SQL 视图的一种通用表示形式,旨在能够被多种查询引擎解释执行。本节介绍如何在 Spark 中使用视图(需要 Spark 3.4 及以上版本,不支持更早的 Spark 版本)。

注意

本节中的所有 SQL 示例均遵循官方 Spark SQL 语法:

创建视图

创建一个不包含任何注释或属性的简单视图:

CREATE VIEW <viewName> AS SELECT * FROM <tableName>

使用 IF NOT EXISTS 可以在视图已存在时避免该 SQL 语句执行失败:

CREATE VIEW IF NOT EXISTS <viewName> AS SELECT * FROM <tableName>

创建一个带注释的视图,其中包含与源表不同的别名列和列注释:

CREATE VIEW <viewName> (ID COMMENT 'Unique ID', ZIP COMMENT 'Zipcode')
    COMMENT 'View Comment'
    AS SELECT id, zip FROM <tableName>

使用属性创建视图

使用 TBLPROPERTIES 创建带有属性的视图:

CREATE VIEW <viewName>
    TBLPROPERTIES ('key1' = 'val1', 'key2' = 'val2')
    AS SELECT * FROM <tableName>

显示视图属性:

SHOW TBLPROPERTIES <viewName>

创建带位置的视图

要指定视图元数据的位置,请使用 TBLPROPERTIES ('location'='完全限定的 URI'):

CREATE VIEW <viewName>
    TBLPROPERTIES ('location' = '/path/to/custom-location')
AS SELECT * FROM <tableName>

视图元数据存储在指定位置,并追加 /metadata,例如 /path/to/custom-location/metadata。

删除视图

删除一个已存在的视图:

DROP VIEW <viewName>

使用 IF EXISTS 可以防止视图不存在时 SQL 语句执行失败:

DROP VIEW IF EXISTS <viewName>

替换视图

使用 CREATE OR REPLACE 可以更新视图的 schema、属性或底层 SQL 语句:

CREATE OR REPLACE VIEW <viewName> (updated_id COMMENT 'updated ID')
    TBLPROPERTIES ('key1' = 'new_val1')
    AS SELECT id FROM <tableName>

设置和移除视图属性

使用 ALTER VIEW ... SET TBLPROPERTIES 设置现有视图的属性:

ALTER VIEW <viewName> SET TBLPROPERTIES ('key1' = 'val1', 'key2' = 'val2')

使用 ALTER VIEW ... UNSET TBLPROPERTIES 从现有视图中移除属性:

ALTER VIEW <viewName> UNSET TBLPROPERTIES ('key1', 'key2')

显示可用视图

列出当前命名空间(通过 USE <namespace> 设置)中的所有视图:

SHOW VIEWS

使用以下变体之一,列出已定义的 catalog 和/或命名空间中的所有可用视图:

SHOW VIEWS IN <catalog>
SHOW VIEWS IN <namespace>
SHOW VIEWS IN <catalog>.<namespace>

显示视图的 CREATE 语句

显示视图的 CREATE 语句:

SHOW CREATE TABLE <viewName>

显示视图详情

使用 DESCRIBE 显示视图的更多详细信息:

DESCRIBE [EXTENDED] <viewName>

评论

登录后参与评论

正在加载评论…