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

PySpark foreach执行S3文件移动时遇FileNotFoundException求助

问题描述

我有一个仅包含一条记录的PySpark DataFrame:

+-----------------------------------------------------------------------+---------------------------------------------------------------------------------+
|source_file_path                                                       |destination_file_path                                                            |
+-----------------------------------------------------------------------+---------------------------------------------------------------------------------+
|Input-Data/EOBLetter/20240425/MONTHLY_EOB_69L7X7F53_20231227212740.json|Input-Data/EOBLetter/Processed/20240425/MONTHLY_EOB_69L7X7F53_20231227212740.json|
+-----------------------------------------------------------------------+---------------------------------------------------------------------------------+

编写了以下类方法,通过foreach()对DataFrame每行执行S3文件复制并删除操作:

class LoadToS3:
    def __init__(self, bucket_name, output_dir):
        self.bucket_name = bucket_name
        self.output_dir = output_dir
        # self.s3 = boto3.resource('s3')

    @property
    def client(self):
        return boto3.client('s3')

    def move_files_in_s3(self, row):
        copy_source = {
            'Bucket': self.bucket_name,
            'Key': row.source_file_path
        }
        
        self.client.copy(copy_source, self.bucket_name,
                        row.destination_file_path)

        self.client.delete_object(
            Bucket=self.bucket_name,
            Key=row.source_file_path
        )

调用方式:

load_to_s3 = LoadToS3(
                    self.args['SOURCE_S3_BUCKET'], self.args['DESTINATION_S3_PATH'])
s3_move_path_df.foreach(load_to_s3.move_files_in_s3)

调用foreach()时出现错误:

Following error occured: An error occurred while calling o384.isEmpty.
: org.apache.spark.SparkException: Job aborted due to stage failure: Task 0 in stage 1374.0 failed 4 times, most recent failure: Lost task 0.3 in stage 1374.0 (TID 5111) (100.127.190.158 executor 1): java.io.FileNotFoundException: 
File not present on S3

It is possible the underlying files have been updated. You can explicitly invalidate
the cache in Spark by running 'REFRESH TABLE tableName' command in SQL or by
recreating the Dataset/DataFrame involved.
解决方案

1. 验证S3文件实际存在性

先确认source_file_path中的文件确实存在于目标S3桶:

  • 用AWS CLI直接检查:
    aws s3 ls s3://<你的桶名>/Input-Data/EOBLetter/20240425/MONTHLY_EOB_69L7X7F53_20231227212740.json
    
  • 注意S3路径区分大小写,检查DataFrame中的路径是否有拼写错误或大小写不一致。

2. 添加文件存在性检查,避免重试报错

Spark任务失败会自动重试,第一次执行时文件已被删除,重试就会触发找不到文件的错误。修改move_files_in_s3方法,先检查文件是否存在:

def move_files_in_s3(self, row):
    source_key = row.source_file_path
    # 检查文件是否存在
    try:
        self.client.head_object(Bucket=self.bucket_name, Key=source_key)
    except self.client.exceptions.ClientError as e:
        if e.response['Error']['Code'] == '404':
            # 文件已被处理或不存在,直接返回
            return
        else:
            # 其他异常抛出
            raise

    copy_source = {
        'Bucket': self.bucket_name,
        'Key': source_key
    }
    
    self.client.copy(copy_source, self.bucket_name,
                    row.destination_file_path)

    self.client.delete_object(
        Bucket=self.bucket_name,
        Key=source_key
    )

3. 改用foreachPartition优化S3客户端创建

当前client属性每次调用都会生成新的boto3客户端,foreach每行创建一个客户端效率极低,还可能引发连接问题。改用foreachPartition,每个分区创建一次客户端:

class LoadToS3:
    def __init__(self, bucket_name, output_dir):
        self.bucket_name = bucket_name
        self.output_dir = output_dir

    def move_files_in_partition(self, partition):
        # 每个分区初始化一次S3客户端
        s3_client = boto3.client('s3')
        for row in partition:
            source_key = row.source_file_path
            try:
                s3_client.head_object(Bucket=self.bucket_name, Key=source_key)
            except s3_client.exceptions.ClientError as e:
                if e.response['Error']['Code'] == '404':
                    continue
                raise

            copy_source = {
                'Bucket': self.bucket_name,
                'Key': source_key
            }
            
            s3_client.copy(copy_source, self.bucket_name,
                          row.destination_file_path)

            s3_client.delete_object(
                Bucket=self.bucket_name,
                Key=source_key
            )

# 调用方式
load_to_s3 = LoadToS3(self.args['SOURCE_S3_BUCKET'], self.args['DESTINATION_S3_PATH'])
s3_move_path_df.foreachPartition(load_to_s3.move_files_in_partition)

4. 临时禁用Spark任务重试

如果确定文件只会被处理一次,可以设置任务最大失败次数为0,避免重试:

from pyspark import SparkContext

sc = SparkContext.getOrCreate()
sc.setLocalProperty("spark.task.maxFailures", "0")

# 执行文件移动操作
s3_move_path_df.foreach(load_to_s3.move_files_in_s3)

注意:此方案仅适用于能确保任务不会因其他原因失败的场景,否则任务会直接终止。

5. 刷新Spark缓存或重建DataFrame

如果DataFrame来自缓存或过时的表数据,可能存在路径失效的情况:

  • 若从表读取,执行SQL刷新:
    REFRESH TABLE your_table_name;
    
  • 若通过读取文件创建DataFrame,重新执行读取逻辑,不要使用缓存版本。

内容的提问来源于stack exchange,提问作者Suraj Pandey

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 05:54:55