【发布时间】:2019-08-04 06:11:33
【问题描述】:
我有一个在 AWS EMR 上运行的 Spark Structured Streaming 任务,它本质上是在一分钟时间窗口内连接两个输入流。输入流有 1 分钟的水印。我不做任何聚合。我使用forEachBatch 和foreachPartition 每个批次“手动”将结果写入S3,将数据转换为字符串并写入S3。
我想运行它很长一段时间,即“永远”,但不幸的是,Spark 会慢慢填满我集群上的 HDFS 存储并最终因此而死。
似乎有两种类型的数据在积累。登录/var 和.delta,.snapshot 文件/mnt/tmp/.../。当我使用 CTRL+C 终止任务(或者在使用 yarn 和 yarn application kill 的情况下)时,它们也不会被删除,我必须手动删除它们。
我使用spark-submit 运行我的任务。我尝试添加
--conf spark.streaming.ui.retainedBatches=100 \
--conf spark.streaming.stopGracefullyOnShutdown=true \
--conf spark.cleaner.referenceTracking.cleanCheckpoints=true \
--conf spark.cleaner.periodicGC.interval=15min \
--conf spark.rdd.compress=true
没有效果。当我添加--master yarn 时,存储临时文件的路径会发生一些变化,但它们随着时间的推移而积累的问题仍然存在。添加--deploy-mode cluster 似乎会使问题变得更糟,因为似乎要写入更多数据。
我的代码中曾经有一个Trigger.ProcessingTime("15 seconds),但当我读到如果触发时间与计算时间相比太短,Spark 可能无法自行清理时将其删除。这似乎有点帮助,HDFS 填充较慢,但临时文件仍在堆积。
如果我不加入这两个流,而只是 select 和 union 将结果写入 S3,则不会发生 cruft int /mnt/tmp 的累积。会不会是我的集群对于输入数据来说太小了?
我想了解 Spark 为什么要编写这些临时文件,以及如何限制它们占用的空间。我也想知道如何限制日志占用的空间量。
【问题讨论】:
标签: apache-spark hdfs amazon-emr spark-structured-streaming