EMR 5.33.1+PySpark 2.4.7迭代写入S3已有目录报错求助
解决方案
1. 使用追加模式写入(推荐)
Spark的DataFrameWriter默认采用error模式(路径已存在时直接抛出异常),只需在写入时指定mode("append"),即可将每次转换后的数据集追加到同一S3分区目录下,不会触发路径已存在的错误。修改后的代码如下:
output_path = 'bucket_uri/folder1/date=20220101' for i in range(0, 100, 10): pdf = spark.read.parquet(file_list[i:i+10]) # ... 执行数据转换操作 ... pdf_transformed.write.mode("append").parquet(output_path)
该方式会在目标目录下新增Parquet文件,所有转换结果最终都保存在指定的分区目录中,完全符合你的需求。
2. 覆盖模式写入(按需选择)
如果你的业务场景需要每次循环都替换目标目录的原有数据,可以使用mode("overwrite"),但该模式会先清空目录下所有已有文件,再写入当前批次的数据,适合全量刷新的场景:
pdf_transformed.write.mode("overwrite").parquet(output_path)
3. 合并DataFrame后一次性写入(小数据量场景)
若处理的数据集规模较小,可以先将每次转换后的DataFrame收集起来,合并为一个大DataFrame后一次性写入,减少IO操作次数:
output_path = 'bucket_uri/folder1/date=20220101' df_list = [] for i in range(0, 100, 10): pdf = spark.read.parquet(file_list[i:i+10]) # ... 执行数据转换操作 ... df_list.append(pdf_transformed) # 初始化空DataFrame并合并所有批次数据 final_df = spark.createDataFrame([], pdf_transformed.schema) for df in df_list: final_df = final_df.union(df) final_df.write.parquet(output_path)
注意:该方式在数据量较大时可能导致Driver内存不足,仅推荐用于中小规模数据集。
内容的提问来源于stack exchange,提问作者haneulkim
相关产品推荐
相关产品推荐

