结构化流处理
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 存储过程。
评论
登录后参与评论
KnowForge