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

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路径下的文件变化:

  1. 第一次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
  1. 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
  1. 执行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手动删除目标路径下的所有对象,但这显然不是最优解,不符合最佳实践。

我的疑问

  1. 有没有正确的方法可以让PySpark在S3上使用overwrite模式时,确保完全覆盖旧数据,避免残留?
  2. 我对PySpark overwrite模式处理S3路径的逻辑理解是否正确?如果是,这种设计是有意为之吗?背后的原因是什么?

备注:内容来源于stack exchange,提问作者emu

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.23 09:54:52