【发布时间】: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