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

Spark处理7z文件时提前删源文件引发FileNotFoundException求助

问题根源

Spark采用惰性求值机制:你在循环中创建的DataFrame仅记录了计算逻辑,并未实际执行数据读取和处理。当最后调用toPandas()时,Spark才会回溯整个执行链尝试读取原始文件,但此时你早已用rm删除了解压后的文件,因此抛出文件不存在的错误。

解决方案

方案1:强制触发计算并持久化结果

在每次循环处理完数据后,通过cache() + count()将筛选后的DataFrame数据持久化到内存(或磁盘),后续union操作将直接使用持久化的数据,不再依赖原始解压文件。修改循环内的代码:

if not df.isEmpty():
    # 持久化数据到内存(若内存不足可改用persist(StorageLevel.DISK_ONLY))
    df = df.cache()
    # 触发实际计算,将数据加载到存储层
    df.count()
    accumulated_df = accumulated_df.union(df)

方案2:分批写入临时存储再合并

如果内存不足以容纳所有数据,可将每个批次的结果写入临时Parquet文件,循环结束后统一读取合并:

def main():
    # ... 其他代码 ...
    import shutil
    temp_dir = "/tmp/weibo_temp_data"
    # 清空临时目录(若存在)
    if os.path.exists(temp_dir):
        shutil.rmtree(temp_dir)
    os.makedirs(temp_dir)

    for file in os.listdir("/home/data/Sources/WEIBO/"):
        print(file)
        if file[-2:] == "7z":
            err = os.system("7z x /home/data/Sources/WEIBO/" + file)
            if err:
                errtime.append(err)
                continue
            df = search_weibo(spark, file[:-3], keywords1, keywords2, keywords3)
            df.show()
            if not df.isEmpty():
                # 写入临时Parquet文件
                df.write.mode("append").parquet(temp_dir)
            os.system("rm " + file[:-3])

    # 循环结束后读取所有临时文件合并
    accumulated_df = spark.read.parquet(temp_dir)
    out = accumulated_df.toPandas()
    out.to_excel("~/ans/ans.xlsx", index=False)
    # 清理临时目录
    shutil.rmtree(temp_dir)

修复代码中的变量名错误

你的search_weibo函数存在变量名错误:读取文件后赋值给f,但后续逻辑用的是df,会导致NameError,需修正:

try:
    df = sparksession.read.schema(custom_schema).format("text").load(filename)  # 将f改为df
except Exception as e:
    print("Cannot load data from " + filename)
    return None

内容的提问来源于stack exchange,提问作者Jade Chen

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 00:02:31