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

PySpark跨权限不同S3桶迁移Parquet数据遇403错误求助

问题:跨S3桶Parquet数据迁移的403权限错误

问题背景

使用EC2实例组成的Spark集群,通过PySpark将Parquet格式数据从源S3桶迁移至目标S3桶,两个桶配置了不同的IAM角色与桶策略。通过修改Hadoop全局配置设置AWS访问密钥,但切换目标桶凭证后,执行df.show()或写入操作时触发403权限拒绝错误。

错误信息

Error occurred: An error occurred while calling o56.showString.
: org.apache.spark.SparkException: Job aborted due to stage failure: Task 0 in stage 1.0 failed 8 times, most recent failure: Lost
task 0.7 in stage 1.0 (TID 8) (172.12.2.153 executor 0):
java.nio.file.AccessDeniedException:
s3a://wayfinder-doceree-s3-customer-data/export/mx_submits_2024_06/part-00000-tid-1074794359726202293-9f504b1d-da56-4031-963b-be9f22348eb4-141340-1.c000.snappy.parquet:
getFileStatus on
s3a://wayfinder-doceree-s3-customer-data/export/mx_submits_2024_06/part-00000-tid-1074794359726202293-9f504b1d-da56-4031-963b-be9f22348eb4-141340-1.c000.snappy.parquet:
com.amazonaws.services.s3.model.AmazonS3Exception: Forbidden (Service:
Amazon S3; Status Code: 403; Error Code: 403 Forbidden; Request ID:
6C07ME6ZAV2B4XSA; S3 Extended Request ID:
Dx8EtSGjnYMl0Ld6kwSs9L9CMk0sdrDkzzdCSsXaG2KXk1uhC6iAIkly0mBCmB6rehqSuat0RlR0WHjPQlFkQQ==;
Proxy: null), S3 Extended Request ID:
Dx8EtSGjnYMl0Ld6kwSs9L9CMk0sdrDkzzdCSsXaG2KXk1uhC6iAIkly0mBCmB6rehqSuat0RlR0WHjPQlFkQQ==:403
Forbidden
at org.apache.hadoop.fs.s3a.S3AUtils.translateException(S3AUtils.java:255).

代码逻辑与问题现象

代码流程:先设置源桶凭证,成功读取Parquet文件到DataFrame,可正常执行df.count()、df.show();切换到目标桶凭证后,执行df.show()或写入操作触发上述错误。

完整代码:

def set_s3_credentials(access_key, secret_key):
    hadoop_conf = spark.sparkContext._jsc.hadoopConfiguration()
    hadoop_conf.set("fs.s3a.access.key", access_key)
    hadoop_conf.set("fs.s3a.secret.key", secret_key)

def copy_parquet_file():
    try:
        # Set and log source AWS credentials
        set_s3_credentials(source_aws_access_key_id, source_aws_secret_access_key)
        logging.info(f"Set source AWS credentials for bucket: '{source_bucket_name}'")
        # update_spark_conf(source_s3_conf)

        logging.info(f"Starting to copy file from bucket '{source_bucket_name}' key '{source_key}' to bucket '{destination_bucket_name}' key '{destination_key}'")
        hadoop_conf = spark.sparkContext._jsc.hadoopConfiguration()
        logging.info(f'Source bucket access_key: {hadoop_conf.get("fs.s3a.access.key")}')
        logging.info(f'Source bucket secret_key: {hadoop_conf.get("fs.s3a.secret.key")}')   

        # Read parquet file from source bucket using PySpark
        source_parquet_path = f"s3a://{source_bucket_name}/{source_key}"
        df = spark.read.parquet(source_parquet_path, header=True, multiLine=True, quote="\"", escape="\"")
        df = df.limit(100) 
        # df.printSchema()
        logging.info("Read source parquet file successfully.")

        # Set and log destination AWS credentials
        set_s3_credentials(dest_aws_access_key_id, dest_aws_secret_access_key)
        # update_spark_conf(destination_s3_conf)
        logging.info(f"Set destination AWS credentials for bucket: '{destination_bucket_name}'")
        
        logging.info(f'Destination bucket access_key: {hadoop_conf.get("fs.s3a.access.key")}')
        logging.info(f'Destination bucket secret_key: {hadoop_conf.get("fs.s3a.secret.key")}')

        # Write dataframe to destination bucket using PySpark
        destination_parquet_path = f"s3a://{destination_bucket_name}/{destination_key}"
        time.sleep(10)
        df.show()
        df.write.parquet(destination_parquet_path, mode='overwrite')
        logging.info(f'Copied {source_key} to {destination_key} successfully.')

    except Exception as e:
        logging.error(f"Error occurred: {e}") 

解决方法

1. 分桶配置凭证(核心解决方案)

Spark DataFrame采用懒加载机制,全局替换Hadoop配置后,executor会用目标桶凭证访问源S3桶(DataFrame依赖仍指向源路径),导致权限不足。需为每个桶单独配置凭证:

修改凭证设置函数:

def set_s3_bucket_credentials(bucket_name, access_key, secret_key):
    hadoop_conf = spark.sparkContext._jsc.hadoopConfiguration()
    # 为指定桶单独配置凭证,不覆盖全局设置
    hadoop_conf.set(f"fs.s3a.bucket.{bucket_name}.access.key", access_key)
    hadoop_conf.set(f"fs.s3a.bucket.{bucket_name}.secret.key", secret_key)

使用方式:

# 初始化时同时设置两个桶的凭证
set_s3_bucket_credentials(source_bucket_name, source_aws_access_key_id, source_aws_secret_access_key)
set_s3_bucket_credentials(destination_bucket_name, dest_aws_access_key_id, dest_aws_secret_access_key)

# 直接读写无需切换凭证
df = spark.read.parquet(source_parquet_path)
df.write.parquet(destination_parquet_path, mode='overwrite')

2. 提前物化DataFrame(备选方案)

如果无法使用分桶配置,可在切换凭证前将DataFrame数据加载到内存,切断与源S3的依赖:

# 读取后立即触发执行,将数据缓存到内存
df = spark.read.parquet(source_parquet_path).limit(100)
df.cache()
df.count()  # 强制执行,加载数据到内存

# 切换凭证
set_s3_credentials(dest_aws_access_key_id, dest_aws_secret_access_key)

# 此时操作不再访问源桶
df.show()
df.write.parquet(destination_parquet_path, mode='overwrite')

注意:数据量较大时缓存可能导致内存不足,需根据实际情况调整。

3. 验证桶策略与IAM权限

  • 源桶凭证需拥有s3:GetObject、s3:ListBucket权限
  • 目标桶凭证需拥有s3:PutObject、s3:ListBucket、s3:DeleteObject(overwrite模式时)权限
  • 确认桶策略允许对应IAM角色/用户访问

4. 使用临时IAM角色凭证(生产环境推荐)

避免硬编码密钥,通过AssumeRole获取临时凭证:

import boto3

def get_temp_credentials(role_arn):
    sts_client = boto3.client('sts')
    response = sts_client.assume_role(RoleArn=role_arn, RoleSessionName="SparkS3Session")
    return response['Credentials']

# 获取源桶临时凭证
source_creds = get_temp_credentials(source_role_arn)
set_s3_bucket_credentials(source_bucket_name, source_creds['AccessKeyId'], source_creds['SecretAccessKey'])

# 获取目标桶临时凭证
dest_creds = get_temp_credentials(dest_role_arn)
set_s3_bucket_credentials(destination_bucket_name, dest_creds['AccessKeyId'], dest_creds['SecretAccessKey'])

内容的提问来源于stack exchange,提问作者Rahul Sheoran

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 08:35:56