Apache Spark

结构化流处理

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

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

Spark Structured Streaming

Iceberg 使用 Apache Spark 的 DataSourceV2 API 来实现数据源与目录(catalog)。Spark DSv2 是一个仍在演进的 API,其在不同 Spark 版本中的支持程度各不相同。

流式读取

Iceberg 支持在 Spark 结构化流作业中从某个历史时间戳开始处理增量数据:

val df = spark.readStream
    .format("iceberg")
    .option("stream-from-timestamp", Long.toString(streamStartTimestamp))
    .load("database.table_name")

警告

Iceberg 仅支持从 append 快照读取数据。overwrite 快照无法被处理,默认情况下会引发异常。可以通过设置 streaming-skip-overwrite-snapshots=true 来忽略 overwrite 快照。类似地,delete 快照默认也会引发异常,可以通过设置 streaming-skip-delete-snapshots=true 来忽略 delete 快照。

限制输入速率

为了控制 DataFrame API 中微批次的大小,Iceberg 支持两个读取选项:

  • streaming-max-files-per-micro-batch 每个微批次中要处理的最大文件数。
  • streaming-max-rows-per-micro-batch 每个微批次中要处理的行数的"软上限"。一个批次总是会包含下一个待处理数据文件中的所有行,但如果包含额外的文件会导致超过该软上限,则不会再包含其他文件。

如果同时设置了这两个选项,微批次的大小将以最先达到限制的那个选项为准。

// Read a hard limit of 1 file per micro-batch
val df = spark.readStream
    .format("iceberg")
    .option("streaming-max-files-per-micro-batch", "1")
    .load("database.table_name")
// Read files until the number of included rows >= 1000 per micro-batch
val df = spark.readStream
    .format("iceberg")
    .option("streaming-max-rows-per-micro-batch", "1000")
    .load("database.table_name")

信息

注意:除了对使用默认触发器(即 Trigger.ProcessingTime)的查询限制微批大小外,速率限制选项还可应用于使用 Trigger.AvailableNow 的查询,以便将所有可用源数据的一次性处理拆分为多个微批,从而提升查询的可扩展性。当使用已弃用的 Trigger.Once 触发器时,速率限制选项将被忽略。

异步微批规划

用户可以通过将 async-micro-batch-planning-enabled 设置为 true 来启用异步微批规划。启用该选项后,Iceberg 会在规划后续微批的同时,并行开始处理当前微批。这有助于减少微批之间的空闲时间,从而提高查询吞吐量。用户应权衡其利弊,其中包括更高的内存占用和增加的快照检测延迟。

用户还可以设置其他选项来控制异步微批规划的行为,这些选项见 Spark 配置。

流式写入

要将流式查询中的值写入 Iceberg 表,请使用 DataStreamWriter:

data.writeStream
    .format("iceberg")
    .outputMode("append")
    .trigger(Trigger.ProcessingTime(1, TimeUnit.MINUTES))
    .option("checkpointLocation", checkpointPath)
    .toTable("database.table_name")

对于基于目录的 Hadoop catalog:

data.writeStream
    .format("iceberg")
    .outputMode("append")
    .trigger(Trigger.ProcessingTime(1, TimeUnit.MINUTES))
    .option("path", "hdfs://nn:8020/path/to/table")
    .option("checkpointLocation", checkpointPath)
    .start()

Iceberg 支持 append 和 complete 输出模式:

  • append:将每个微批的行追加到表中
  • complete:在每个微批中替换表的全部内容

在启动流式查询之前,请确保已经创建了该表。有关如何创建 Iceberg 表,请参阅 SQL create table 文档。

Iceberg 不支持实验性的连续处理,因为它没有提供"提交"输出的接口。

分区表

Iceberg 要求在写入数据之前,按任务对数据进行分区排序。在 Spark 中,针对分区表,任务会按 Spark 分区进行划分。对于批处理查询,建议你显式进行排序以满足该要求(参见此处),但这种方式会带来额外的延迟,因为对于流式工作负载而言,重新分区和排序都属于开销较大的操作。为避免额外的延迟,你可以启用 fanout writer 来消除这一要求。

data.writeStream
    .format("iceberg")
    .outputMode("append")
    .trigger(Trigger.ProcessingTime(1, TimeUnit.MINUTES))
    .option("fanout-enabled", "true")
    .option("checkpointLocation", checkpointPath)
    .toTable("database.table_name")

扇出写入器(fanout writer)会按分区值打开文件,并在写入任务结束前不会关闭这些文件。请避免在批量写入中使用扇出写入器,因为对于批处理工作负载而言,对输出行进行显式排序的开销很低。

流式表的维护

流式写入可能会快速创建新的表版本,从而产生大量用于跟踪这些版本的表元数据。建议通过调整提交速率、过期旧快照以及自动清理元数据文件来维护元数据。

调整提交速率

提交频率过高会产生大量数据文件、manifest 和快照,进而带来额外的维护工作。建议触发间隔至少为 1 分钟,并根据需要适当延长该间隔。

Structured Streaming 编程指南中的 triggers(触发器)部分说明了如何配置该间隔。

过期旧快照

写入表的每个批次都会产生一个新的快照。Iceberg 会在表元数据中跟踪这些快照,直到它们过期为止。在频繁提交的情况下,快照会迅速累积,因此强烈建议对由流式查询写入的表定期进行维护。快照过期是指移除不再需要的元数据及相关数据文件的过程。默认情况下,该过程会将五天前的快照标记为过期。

压缩数据文件

流式处理写入的数据量通常较小,这可能导致表元数据中记录大量小文件。将小文件压缩为大文件可以减少表所需的元数据量,并提高查询效率。Iceberg 和 Spark 提供了 rewrite_data_files 存储过程。

重写 manifest

为了优化流式工作负载的写入延迟,Iceberg 可以使用"快速"追加(fast append)方式写入新快照,这种方式不会自动合并 manifest。这可能导致产生大量小的 manifest 文件。Iceberg 可以重写 manifest 文件以提升查询性能。Iceberg 和 Spark 提供了 rewrite_manifests 存储过程。

评论

登录后参与评论

正在加载评论…