【问题标题】:How to manage HDFS memory with Structured Streaming Checkpoints如何使用结构化流检查点管理 HDFS 内存
【发布时间】:2019-06-01 21:21:20
【问题描述】:

我有一个长期运行的结构化流作业,它消耗多个 Kafka 主题并在一个滑动窗口上聚合。我需要了解如何在 HDFS 中管理/清理检查点。

作业运行良好,我能够从失败的步骤中恢复而不会丢失任何数据,但是,我可以看到 HDFS 利用率每天都在增加。我找不到任何关于 Spark 如何管理/清理检查点的文档。以前检查点存储在 s3 上,但结果证明这是非常昂贵的,因为需要读取/写入大量小文件。

query = formatted_stream.writeStream \
                        .format("kafka") \
                        .outputMode(output_mode) \
                        .option("kafka.bootstrap.servers", bootstrap_servers) \
                        .option("checkpointLocation", "hdfs:///path_to_checkpoints") \
                        .start()

据我了解,检查点应自动清理;几天后,我看到我的 HDFS 利用率呈线性增长。如何确保检查点得到管理且 HDFS 不会耗尽空间?

Spark Structured Streaming Checkpoint Cleanup 接受的答案表明结构化流应该处理这个问题,而不是如何或如何配置它。

【问题讨论】:

标签: apache-spark hdfs spark-structured-streaming


【解决方案1】:

正如您在the code for Checkpoint.scala 中看到的那样,检查点机制会保留最后 10 个检查点数据,但这在几天内应该不是问题。

通常的原因是,您在磁盘上持久化的 RDD 也随着时间线性增长。这可能是由于某些您不关心的 RDD 会被持久化。

您需要确保在使用结构化流时没有需要持久化的增长的 RDD。例如,如果您想计算数据集列上不同元素的精确计数,您需要知道完整的输入数据(这意味着如果您每批次有恒定的数据流入,则持久数据会随时间线性增加)。相反,如果您可以使用近似计数,则可以使用 HyperLogLog++ 等算法,这通常需要更少的内存来权衡精度。

请记住,如果您使用的是 Spark SQL,您可能需要进一步检查优化后的查询变成了什么,因为这可能与 Catalyst 如何优化您的查询有关。如果您不是,那么如果您这样做了,Catalyst 可能会为您优化查询。

在任何情况下,进一步思考:如果检查点的使用随着时间的推移而增加,这应该反映在您的流作业也随着时间线性消耗更多的 RAM,因为检查点只是 Spark 上下文的序列化(加上常数-size 元数据)。如果是这种情况,请查看 SO 以了解相关问题,例如 why does memory usage of Spark Worker increase with time?

此外,请注意您在哪些 RDD 上调用 .persist()(以及哪个缓存级别,以便您可以元数据到磁盘 RDD 并且一次仅将它们部分加载到 Spark 上下文中)。

【讨论】:

  • 感谢您的洞察力。关于持久化数据,我对窗口和 id 字段(min、max、first、sum、count、avg、collect_set 和 approx_count_distinct)执行聚合。我曾假设这不会持续超过窗口+水印?同样,我在带水印流的列子集上使用 drop_duplicates()。这会导致检查点被持久化吗?
  • 我对此进行了进一步调查,发现根本原因是调用了 dropDuplicates()。我误解了有关水印的文档,导致快照随着时间的推移逐渐增长。再次感谢您的洞察力。