PySpark跨权限不同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

