PySpark写入S3时mode('overwrite')未正确覆盖数据的问题咨询及解决方案求助
PySpark写入S3时mode('overwrite')未正确覆盖数据的问题咨询及解决方案求助
大家好,我在项目中遇到了PySpark写入S3时mode('overwrite')无法正确覆盖数据的问题,想请教社区的大佬们帮忙分析和解决。
问题背景
我们项目使用PySpark处理数据,需要将结果存储到Amazon S3中,但发现当目标路径下已经存在文件时,使用pyspark.sql.DataFrame.write配合mode="overwrite"无法正确覆盖S3中的数据,旧数据会残留下来。
复现步骤
下面是完整的复现代码和操作流程:
0. 初始化环境和配置
import pandas as pd import awswrangler as wr from pyspark.sql import SparkSession spark = SparkSession.builder \ .config('spark.hadoop.fs.s3a.aws.credentials.provider', 'com.amazonaws.auth.profile.ProfileCredentialsProvider') \ .config('spark.hadoop.fs.s3a.bucket.all.committer.magic.enabled', 'true') \ .getOrCreate() # PySpark使用的输出路径,采用s3a://协议 output_url = f's3a://{MY_BUCKET}/test.csv' # awswrangler使用的同路径,转换为s3://协议 wr_output_url = output_url.replace('s3a:', 's3:')
1. 第一次用PySpark写入数据
df = spark.createDataFrame([{'Key': 'OldFoo'}, {'Key': 'OldBar'}], ['Key']) df.write.csv(output_url)
2. 用awswrangler写入同路径
wr.s3.to_csv(pd.DataFrame([{'SomeKey': 'SomeValue'} ]), wr_output_url)
3. 用PySpark的overwrite模式写入新数据
df = spark.createDataFrame([{'Key': 'Foo'}, {'Key': 'Bar'}], ['Key']) df.write.mode('overwrite').csv(output_url)
4. 读取数据验证结果
spark.read.csv(output_url).show()
预期输出
只显示最后一次写入的两条数据:
+------+ | _c0| +------+ | Foo| | Bar| +------+
实际输出
旧数据和新数据同时存在:
+------+ | _c0| +------+ |OldFoo| |OldBar| | Foo| | Bar| +------+
现象分析
我查看了每一步操作后S3路径下的文件变化:
- 第一次PySpark写入后,S3中的文件:
s3://<MY_BUCKET>/test.csv/_SUCCESS s3://<MY_BUCKET>/test.csv/part-00000-52358825-0cf7-4609-81b1-2819d4205d85-c000.csv s3://<MY_BUCKET>/test.csv/part-00001-52358825-0cf7-4609-81b1-2819d4205d85-c000.csv s3://<MY_BUCKET>/test.csv/part-00003-52358825-0cf7-4609-81b1-2819d4205d85-c000.csv
- awswrangler写入后,多了一个直接的
test.csv文件:
s3://<MY_BUCKET>/test.csv s3://<MY_BUCKET>/test.csv/_SUCCESS s3://<MY_BUCKET>/test.csv/part-00000-52358825-0cf7-4609-81b1-2819d4205d85-c000.csv s3://<MY_BUCKET>/test.csv/part-00001-52358825-0cf7-4609-81b1-2819d4205d85-c000.csv s3://<MY_BUCKET>/test.csv/part-00003-52358825-0cf7-4609-81b1-2819d4205d85-c000.csv
- 执行PySpark overwrite后,发现旧的part文件依然存在,只删除了直接的
test.csv文件:
s3://<MY_BUCKET>/test.csv/_SUCCESS s3://<MY_BUCKET>/test.csv/part-00000-503a773b-4f7d-4089-9bce-f87bf56eb3df-c000.csv s3://<MY_BUCKET>/test.csv/part-00000-52358825-0cf7-4609-81b1-2819d4205d85-c000.csv s3://<MY_BUCKET>/test.csv/part-00001-503a773b-4f7d-4089-9bce-f87bf56eb3df-c000.csv s3://<MY_BUCKET>/test.csv/part-00001-52358825-0cf7-4609-81b1-2819d4205d85-c000.csv s3://<MY_BUCKET>/test.csv/part-00003-503a773b-4f7d-4089-9bce-f87bf56eb3df-c000.csv s3://<MY_BUCKET>/test.csv/part-00003-52358825-0cf7-4609-81b1-2819d4205d85-c000.csv
看起来问题出在:当S3路径下同时存在一个同名文件(test.csv)和以该文件名作为前缀的目录结构(test.csv/xxx)时,PySpark的overwrite模式只会删除那个单独的文件,而不会清理前缀目录下的旧文件。而如果之前只有前缀目录没有单独文件时,overwrite是正常工作的。
这应该和S3是键值存储而非传统文件系统的特性有关,传统文件系统不允许文件和同名目录共存,但S3支持这种情况。
已尝试的研究和临时方案
- 查看Hadoop相关逻辑:我推测PySpark底层使用的Hadoop S3A客户端在处理删除时,逻辑是如果是目录就递归删除内容,如果是单个文件就只删文件,不会处理“文件+前缀目录”共存的情况。
- 临时 workaround:每次写入前用awswrangler手动删除目标路径下的所有对象,但这显然不是最优解,不符合最佳实践。
我的疑问
- 有没有正确的方法可以让PySpark在S3上使用
overwrite模式时,确保完全覆盖旧数据,避免残留? - 我对PySpark
overwrite模式处理S3路径的逻辑理解是否正确?如果是,这种设计是有意为之吗?背后的原因是什么?
备注:内容来源于stack exchange,提问作者emu
相关产品推荐
相关产品推荐

