【发布时间】:2023-04-02 00:46:01
【问题描述】:
我有一个流作业,它旨在通过使用mapWithState 的单个步骤连续运行,因此需要配置检查点。我使用本地目录设置它,因为在这个阶段它只在单个节点上运行。
我观察到检查点目录快速且持续增长。在几天的时间里,它会增长到超过一百万个文件并耗尽磁盘上的 inode。
问题:
- 这是预期行为吗?
- 假设不是,我如何隔离可能导致快照不被修剪的原因?
【问题讨论】:
标签: apache-spark spark-streaming
我有一个流作业,它旨在通过使用mapWithState 的单个步骤连续运行,因此需要配置检查点。我使用本地目录设置它,因为在这个阶段它只在单个节点上运行。
我观察到检查点目录快速且持续增长。在几天的时间里,它会增长到超过一百万个文件并耗尽磁盘上的 inode。
【问题讨论】:
标签: apache-spark spark-streaming
错误是检查点是由sparkContext.checkpoint(checkpointDir) 而不是sparkStreamingContext.checkpoint(checkpointDir) 启用的。
前者足以让 Spark 运行有状态流,而不是抱怨未启用检查点,但未调用流检查点的适当逻辑,因为 sparkStreamingContext.checkpointDir 为空。
【讨论】: