Spark Structured Streaming v2.4.0 checkpoint目录持续增长:*.tmp.crc文件未自动删除?
这确实是Spark 2.4版本里的一个已知问题——严格来说,它不算Spark结构化流本身的bug,而是Hadoop文件系统客户端与Spark checkpoint清理逻辑的协作漏洞。
问题原因
这些.tmp.crc文件是Hadoop FileSystem在写入临时checkpoint文件时生成的校验文件,用于验证文件完整性。Spark在清理旧批次时,只会删除commits、offsets、state目录下的业务文件(比如你看到的6605、6606这类批次号文件),但没有关联删除对应的临时校验文件,导致这些crc文件不断积累,最终撑大checkpoint目录。即使设置了spark.sql.streaming.minBatchesToRetain,也只影响业务文件的保留数量,对这些隐藏的crc文件无效。
解决方案
1. 升级Spark版本(最彻底的方案)
Spark 3.0及后续版本已经修复了这个问题,优化了checkpoint清理逻辑,会自动关联删除对应的.tmp.crc文件。如果你的业务场景允许升级,这是一劳永逸的解决办法。
2. 编写定时清理脚本(临时过渡方案)
如果暂时无法升级,可以写一个Shell脚本,定期清理超过保留批次范围的crc文件。脚本逻辑可以结合你设置的spark.sql.streaming.minBatchesToRetain参数,只保留与当前活跃批次对应的crc文件:
#!/bin/bash # 配置你的checkpoint路径和保留批次数量 CHECKPOINT_DIR="/your/checkpoint/path" MIN_BATCHES_RETAIN=10 # 与spark.sql.streaming.minBatchesToRetain保持一致 # 获取commits目录下最新的批次号 MAX_BATCH=$(ls -1 "${CHECKPOINT_DIR}/commits" | sort -n | tail -1) # 计算需要保留的最小批次号 MIN_BATCH=$((MAX_BATCH - MIN_BATCHES_RETAIN)) # 清理指定目录下的过期crc文件 clean_dir() { local dir_path=$1 local batch_extractor=$2 cd "${dir_path}" for file in .*.tmp.crc; do # 跳过目录本身和上级目录 if [[ "$file" == "." || "$file" == ".." ]]; then continue fi # 提取批次号 BATCH_NUM=$($batch_extractor "$file") # 删除小于最小保留批次的文件 if [[ $BATCH_NUM -lt $MIN_BATCH ]]; then rm -f "$file" echo "Deleted: ${dir_path}/${file}" fi done } # 清理commits目录(格式:..100.xxxx.tmp.crc) clean_dir "${CHECKPOINT_DIR}/commits" "echo \$1 | awk -F '.' '{print \$3}'" # 清理offsets目录(格式同commits) clean_dir "${CHECKPOINT_DIR}/offsets" "echo \$1 | awk -F '.' '{print \$3}'" # 清理state目录(格式:..00000000000000001234.tmp.crc) clean_dir "${CHECKPOINT_DIR}/state" "echo \$1 | sed -E 's/.*\\.([0-9]+)\\.tmp\\.crc/\\1/'"
你可以把这个脚本加入crontab,设置每天低峰期自动运行(比如凌晨2点):
0 2 * * * /path/to/your/clean_crc_files.sh >> /var/log/clean_crc.log 2>&1
注意:运行脚本时要确保Spark任务没有在写入该checkpoint文件,避免出现文件冲突。
3. 调整Hadoop配置(不推荐)
可以尝试修改Hadoop的配置参数,比如设置fs.hdfs.impl.disable.cache=true来禁用HDFS客户端缓存,或者调整校验文件的生成策略,但这种方法可能会影响其他依赖Hadoop的操作性能,而且不一定能完全解决问题,所以不建议作为首选方案。
总结
如果条件允许,优先升级到Spark 3.x版本;如果暂时无法升级,定时脚本是最稳妥的过渡方案。
内容的提问来源于stack exchange,提问作者Nebo Topo

