Apache Spark

写入

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

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

Spark 写入

要在 Spark 中使用 Iceberg,首先需要配置 Spark catalog。

某些计划(plan)只有在使用 Iceberg SQL 扩展时才可用。

Iceberg 使用 Apache Spark 的 DataSourceV2 API 来实现数据源和 catalog。Spark DSv2 是一个不断演进的 API,在不同版本的 Spark 中支持程度各不相同:

功能支持Spark说明
SQL insert into✔️⚠ 需要 spark.sql.storeAssignmentPolicy=ANSI(自 Spark 3.0 起为默认值)
SQL merge into✔️⚠ 需要 Iceberg Spark 扩展
SQL insert overwrite✔️⚠ 需要 spark.sql.storeAssignmentPolicy=ANSI(自 Spark 3.0 起为默认值)
SQL delete from✔️⚠ 行级删除需要 Iceberg Spark 扩展
SQL update✔️⚠ 需要 Iceberg Spark 扩展
DataFrame 追加✔️
DataFrame 覆写✔️
DataFrame CTAS 与 RTAS✔️⚠ 需要 DSv2 API
DataFrame merge into✔️⚠ 需要 DSv2 API(Spark 4.0 及更高版本)

使用 SQL 写入

Spark 支持 SQL 的 INSERT INTO、MERGE INTO 和 INSERT OVERWRITE,以及新的 DataFrameWriterV2 API。

INSERT INTO

要向表中追加新数据,请使用 INSERT INTO。

INSERT INTO prod.db.table VALUES (1, 'a'), (2, 'b')
INSERT INTO prod.db.table SELECT ...

MERGE INTO

Spark 支持 MERGE INTO 查询,它可以表达行级更新。

Iceberg 通过在一次 overwrite 提交中重写包含需要更新的行的数据文件来支持 MERGE INTO。

建议使用 MERGE INTO 而不是 INSERT OVERWRITE,因为 Iceberg 只能替换受影响的数据文件,并且动态覆盖所覆盖的数据可能会随着表的分区方式变化而改变。

MERGE INTO 语法

MERGE INTO 使用来自另一个查询(称为 源)的一组更新来更新一个表(称为 目标 表)。目标表中某一行的更新通过 ON 子句找到,该子句类似于连接条件。

MERGE INTO prod.db.target t   -- a target table
USING (SELECT ...) s          -- the source updates
ON t.id = s.id                -- condition to find updates for target rows
WHEN ...                      -- updates

使用 WHEN MATCHED ... THEN ... 列出对目标表中行的更新操作。可以添加多个 MATCHED 子句,并附带条件以确定何时应用相应的匹配,系统会使用第一个匹配的表达式。

WHEN MATCHED AND s.op = 'delete' THEN DELETE
WHEN MATCHED AND t.count IS NULL AND s.op = 'increment' THEN UPDATE SET t.count = 0
WHEN MATCHED AND s.op = 'increment' THEN UPDATE SET t.count = t.count + 1

未匹配的源行(更新)可以插入:

WHEN NOT MATCHED THEN INSERT *

插入还支持附加条件:

WHEN NOT MATCHED AND s.event_time > still_valid_threshold THEN INSERT (id, count) VALUES (s.id, 1)

源数据中只能有一条记录更新目标表中的任意给定行,否则将抛出错误。

Spark 3.5 新增了对 WHEN NOT MATCHED BY SOURCE ... THEN ... 的支持,可用于更新或删除源数据中不存在的行:

WHEN NOT MATCHED BY SOURCE THEN UPDATE SET status = 'invalid'

快照摘要

在 MERGE INTO 提交之后,快照摘要可能包含以下字段。每个值都是非负计数的字符串形式。当值未知时(例如,未由 Spark 报告),该字段会被省略。

信息

仅在 Spark 4.1 及更高版本中可用。

字段说明
spark.merge-into.num-target-rows-copied因未匹配任何操作而被原样复制的目标行数
spark.merge-into.num-target-rows-deleted被删除的目标行数
spark.merge-into.num-target-rows-updated被更新的目标行数
spark.merge-into.num-target-rows-inserted被插入的目标行数
spark.merge-into.num-target-rows-matched-updated由 MATCHED 子句更新的目标行数
spark.merge-into.num-target-rows-matched-deleted由 MATCHED 子句删除的目标行数
spark.merge-into.num-target-rows-not-matched-by-source-updated由 NOT MATCHED BY SOURCE 子句更新的目标行数
spark.merge-into.num-target-rows-not-matched-by-source-deleted由 NOT MATCHED BY SOURCE 子句删除的目标行数

INSERT OVERWRITE

INSERT OVERWRITE 可以用查询结果替换表中的数据。对于 Iceberg 表,覆写是原子操作。

INSERT OVERWRITE 将替换哪些分区,取决于 Spark 的分区覆写模式和表的分区方式。MERGE INTO 只能重写受影响的数据文件,其行为更容易理解,因此建议使用它来代替 INSERT OVERWRITE。

覆写行为

Spark 默认的覆写模式是静态模式,但在写入 Iceberg 表时建议使用动态覆写模式。 静态覆写模式通过将 PARTITION 子句转换为过滤条件来确定要覆写表中的哪些分区,但 PARTITION 子句只能引用表的列。

动态覆写模式通过设置 spark.sql.sources.partitionOverwriteMode=dynamic 进行配置。

为了演示动态覆写和静态覆写的行为,请看以下 DDL 定义的 logs 表:

CREATE TABLE prod.my_app.logs (
    uuid string NOT NULL,
    level string NOT NULL,
    ts timestamp NOT NULL,
    message string)
USING iceberg
PARTITIONED BY (level, hours(ts))

动态覆盖

当 Spark 的覆盖模式为动态时,SELECT 查询产生数据的分区将被替换。

例如,下面的查询会从示例 logs 表中删除重复的日志事件。

INSERT OVERWRITE prod.my_app.logs
SELECT uuid, first(level), first(ts), first(message)
FROM prod.my_app.logs
WHERE cast(ts as date) = '2020-07-01'
GROUP BY uuid

在动态模式下,这将替换掉任何在 SELECT 结果中存在行的分区。由于所有行的日期都被限制在 7 月 1 日,因此只有当天的小时分区会被替换。

静态覆盖

当 Spark 的覆盖模式为静态时,PARTITION 子句会被转换为一个过滤器,用于从表中删除数据。如果省略了 PARTITION 子句,则所有分区都会被替换。

由于上面的查询中没有 PARTITION 子句,在静态模式下运行时会删除表中所有已有的行,但只会写入 7 月 1 日的日志数据。

若只想覆盖已加载的分区,请添加与 SELECT 查询过滤条件相对应的 PARTITION 子句:

INSERT OVERWRITE prod.my_app.logs
PARTITION (level = 'INFO')
SELECT uuid, first(level), first(ts), first(message)
FROM prod.my_app.logs
WHERE level = 'INFO'
GROUP BY uuid

注意,此模式无法像动态分区示例查询那样替代小时级分区,因为 PARTITION 子句只能引用表列,不能引用隐藏分区。

DELETE FROM

Spark 支持 DELETE FROM 查询,用于从表中删除数据。

删除查询接受一个过滤条件,以匹配要删除的行。

DELETE FROM prod.db.table
WHERE ts >= '2020-05-01 00:00:00' and ts < '2020-06-01 00:00:00'

DELETE FROM prod.db.all_events
WHERE session_time < (SELECT min(session_time) FROM prod.db.good_events)

DELETE FROM prod.db.orders AS t1
WHERE EXISTS (SELECT oid FROM prod.db.returned_orders WHERE t1.oid = oid)

如果删除过滤条件匹配的是表的整个分区,Iceberg 将执行仅元数据的删除操作;如果匹配的是表中的单个行,则 Iceberg 只会重写受影响的数据文件。

UPDATE

更新查询接受一个过滤条件,用于指定要更新的行。

UPDATE prod.db.table
SET c1 = 'update_c1', c2 = 'update_c2'
WHERE ts >= '2020-05-01 00:00:00' and ts < '2020-06-01 00:00:00'

UPDATE prod.db.all_events
SET session_time = 0, ignored = true
WHERE session_time < (SELECT min(session_time) FROM prod.db.good_events)

UPDATE prod.db.orders AS t1
SET order_status = 'returned'
WHERE EXISTS (SELECT oid FROM prod.db.returned_orders WHERE t1.oid = oid)

如需基于传入数据进行更复杂的行级更新,请参阅 MERGE INTO 章节。

写入分支

执行写入前,分支必须已存在。如果分支不存在,操作不会自动创建分支。可以使用 Spark DDL 创建分支。

Info

注意:向分支写入时,将使用表的当前 schema 进行校验。

通过 SQL

在操作中提供分支标识符 branch_yourBranch,即可执行分支写入。

通过指定 spark.wap.branch 配置,分支写入也可以作为写入-审计-发布(WAP)工作流的一部分执行。注意,WAP 分支和分支标识符不能同时指定。

-- INSERT (1,' a') (2, 'b') into the audit branch.
INSERT INTO prod.db.table.branch_audit VALUES (1, 'a'), (2, 'b');

-- MERGE INTO audit branch
MERGE INTO prod.db.table.branch_audit t
USING (SELECT ...) s
ON t.id = s.id
WHEN ...

-- UPDATE audit branch
UPDATE prod.db.table.branch_audit AS t1
SET val = 'c'

-- DELETE FROM audit branch
DELETE FROM prod.db.table.branch_audit WHERE id = 2;

-- WAP Branch write
SET spark.wap.branch = audit-branch
INSERT INTO prod.db.table VALUES (3, 'c');

通过 DataFrame

通过 DataFrame 写入分支时,只需在操作中提供分支标识符 branch_yourBranch 即可。

// To insert into `audit` branch
val data: DataFrame = ...
data.writeTo("prod.db.table.branch_audit").append()
// To overwrite `audit` branch
val data: DataFrame = ...
data.writeTo("prod.db.table.branch_audit").overwritePartitions()

使用 DataFrame 写入

Spark 引入了新的 DataFrameWriterV2 API,用于通过数据框写入表。推荐使用 v2 API 的原因如下:

  • 支持 CTAS、RTAS 以及按过滤条件覆盖写入

  • 所有操作都一致地按列名向表中写入列

  • partitionedBy 中支持隐藏分区表达式

  • 覆盖行为是显式的,可以是动态覆盖,也可以由用户提供的过滤条件指定

  • 每个操作的行为都与 SQL 语句相对应

    • df.writeTo(t).create() 等价于 CREATE TABLE AS SELECT
    • df.writeTo(t).replace() 等价于 REPLACE TABLE AS SELECT
    • df.writeTo(t).append() 等价于 INSERT INTO
    • df.writeTo(t).overwritePartitions() 等价于动态 INSERT OVERWRITE

v1 DataFrame 的 write API 仍然受支持,但不推荐使用。

危险

在 Spark 中使用 v1 DataFrame API 写入时,请使用 saveAsTable 或 insertInto 通过 catalog 加载表。使用 format("iceberg") 会加载一个孤立的表引用,它不会自动刷新查询所使用的表。

追加数据

要将数据框追加到 Iceberg 表中,请使用 append:

val data: DataFrame = ...
data.writeTo("prod.db.table").append()

覆写数据

要动态覆写分区,请使用 overwritePartitions():

val data: DataFrame = ...
data.writeTo("prod.db.table").overwritePartitions()

要显式覆盖分区,请使用 overwrite 提供过滤条件:

data.writeTo("prod.db.table").overwrite($"level" === "INFO")

创建表

要运行 CTAS 或 RTAS,请使用 create、replace 或 createOrReplace 操作:

val data: DataFrame = ...
data.writeTo("prod.db.table").create()

如果你已将默认的 Spark catalog(spark_catalog)替换为 Iceberg 的 SparkSessionCatalog,请执行以下操作:

val data: DataFrame = ...
data.writeTo("db.table").using("iceberg").create()

创建与替换操作支持表配置方法,例如 partitionedBy 和 tableProperty:

data.writeTo("prod.db.table")
    .tableProperty("write.format.default", "orc")
    .partitionedBy($"level", days($"ts"))
    .createOrReplace()

Iceberg 表的位置也可以通过 location 表属性来指定:

data.writeTo("prod.db.table")
    .tableProperty("location", "/path/to/location")
    .createOrReplace()

合并数据

Spark 4.0 新增了使用 DataFrameWriterV2 API 执行 MERGE INTO 查询的支持。

MERGE INTO 查询使用来自源的一组更新来更新目标表,在此场景中,源是一个 DataFrame:

val source: DataFrame = ...                               // e.g., read from a table, "source"
source.mergeInto("target", $"source.id" === $"target.id") // second argument is the ON condition
    .whenMatched($"target.id" === 1)                      // argument is the additional condition
    .updateAll()                                          // UPDATE SET *
    .whenMatched($"target.id" === 2)
    .delete()
    .whenNotMatched()
    .insertAll()                                          // INSERT *
    .whenNotMatchedBySource($"target.id" === 3)
    .update(Map("status" -> lit("invalid")))              // set column name(s) to expression(s)
    .merge()

模式合并

在插入或更新时,Iceberg 能够在运行时解析模式不匹配的问题。如果进行了相应配置,Iceberg 将执行以下自动模式演进:

  • 源数据中存在新列,而目标表中没有该列。

    新列会被添加到目标表中。表中已存在的所有行,其该列的值均设置为 NULL

  • 目标表中存在列,而源数据中没有该列。

    插入时该列的值设置为 NULL,更新时则保持原值不变。

必须将目标表的属性 write.spark.accept-any-schema 设置为 true,目标表才能接受任何模式变更。

ALTER TABLE prod.db.sample SET TBLPROPERTIES (
  'write.spark.accept-any-schema'='true'
)

写入方必须启用 mergeSchema 选项。

data.writeTo("prod.db.sample").option("mergeSchema","true").append()

写入分布模式

Iceberg 默认的 Spark 写入器要求每个 Spark 任务中的数据按分区值进行聚类。要求这种分布是为了在写入期间尽量减少保持打开状态的文件句柄数量。从 Iceberg 1.2.0 开始,Iceberg 默认还会请求 Spark 对待写入的数据进行预排序,以符合这种分布。该请求通过表属性 write.distribution-mode 传递给 Spark,其值为 hash。Spark 3.5.0 之前的版本在 CTAS/RTAS 中不会遵守分布模式。

下面我们针对如下示例表来讲解数据写入:

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

要向示例表写入数据,需要按 days(ts), category 对数据进行排序,但这会由默认的 hash 分布方式自动完成。此前这需要手动排序,但现在不再需要了。

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

write.distribution-mode 有 3 种选项

  • none - 这是 Iceberg 之前的默认模式。
    该模式不会请求 Spark 自动执行任何 shuffle 或排序。由于 Spark 不会自动完成任何工作,数据必须手动按分区值排序。数据必须在每个 Spark task 内部排序,或在整个数据集范围内全局排序。全局排序可以最小化输出文件的数量。
    可以使用 Spark 的 write fanout 属性来避免排序,但这会导致所有文件句柄在每个写入 task 完成之前一直保持打开状态。
  • hash - 这是新的默认模式,它请求 Spark 使用基于哈希的交换(exchange)在写入之前对传入的写入数据进行 shuffle。
    实际上,这意味着每一行会根据其分区值进行哈希计算,然后依据该值被放置到对应的 Spark task 中。由于 Spark 的自适应查询规划,task 可能会进一步被拆分或合并。
  • range - 该模式请求 Spark 在写入之前执行基于范围的交换来 shuffle 数据。
    这是一个两阶段的过程,比 hash 模式开销更大。第一阶段会根据分区列和排序列对待写入的数据进行采样。第二阶段利用范围信息将输入数据 shuffle 到各个 Spark task 中。每个 task 获得一个独占的输入数据范围,从而使数据按分区聚类,同时也实现了全局排序。
    虽然这种方式比哈希分布开销更大,但如果查询中使用了排序列,全局有序对读取性能是有益的。如果表创建时指定了排序规则(sort-order),则默认使用该模式。由于 Spark 的自适应查询规划,task 可能会进一步被拆分或合并。

控制文件大小

使用 Spark 将数据写入 Iceberg 时,需要注意 Spark 无法写出比单个 Spark task 更大的文件,且文件不能跨越 Iceberg 分区边界。这意味着,虽然 Iceberg 总会在文件增长到 write.target-file-size-bytes 时滚动文件,但如果 Spark task 本身不够大,这种情况就不会发生。磁盘上创建的文件大小也会远小于 Spark task,因为磁盘上的数据经过了压缩且采用列式存储格式,而 Spark 使用的是未压缩的行式表示。这意味着一个 100 兆字节的 Spark task 即使只写入单个 Iceberg 分区,创建的文件也会远小于 100 兆字节。如果 task 写入多个分区,文件会更小。

要控制哪些数据进入每个 Spark 任务,请使用 write distribution mode,或者手动对数据重新分区。

要调整 Spark 的任务大小,熟悉 Spark 的各项 Adaptive Query Execution(AQE,自适应查询执行)参数非常重要。当 write.distribution-mode 不为 none 时,AQE 会在数据交换(exchange)过程中控制 Spark 任务的合并与拆分,以尝试生成大小为 spark.sql.adaptive.advisoryPartitionSizeInBytes 的任务。这些设置同样会影响用户自行执行的重新分区或排序操作。需要再次注意的是,这是 Spark 内存中行的大小,而非磁盘上经过列式压缩后的大小,因此需要指定一个比目标文件大小更大的值。内存中大小与磁盘上大小的比率取决于数据本身。Spark 后续的工作应能允许 Iceberg 在写入时自动调整该参数,使其与 write.target-file-size-bytes 相匹配。

评论

登录后参与评论

正在加载评论…