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
相关产品推荐
相关产品推荐

