You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Spark Structured Streaming v2.4.0 checkpoint目录持续增长:*.tmp.crc文件未自动删除?

Spark 2.4结构化流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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.12 03:53:49