如何在PySpark中使用两套AWS凭证跨账号读写S3存储桶
跨AWS账号S3数据迁移的凭证冲突问题
问题背景
需要从某AWS账号的S3存储桶读取多份文件,写入另一AWS账号的S3存储桶。
现状问题
在函数中切换AWS凭证后,执行写入操作时出现读取相关错误——Spark执行写入时仍使用新凭证去访问原存储桶,导致权限不足。
问题原因
- Spark的Hadoop配置是全局共享的,修改
spark._jsc.hadoopConfiguration()的凭证会直接覆盖全局设置。 - Spark采用惰性执行机制:
read.parquet操作不会在read_from_bucket1中立即执行,而是等到write_to_bucket2触发Action(如save)时才真正读取数据,此时全局凭证已经被切换为目标账号的,自然无法访问原账号的S3桶。
解决方案
方法1:为不同S3桶配置独立凭证(推荐)
利用Hadoop S3A的按桶配置特性,直接为每个桶单独设置凭证,避免全局覆盖。配置格式为fs.s3a.bucket.<bucket-name>.access.key和fs.s3a.bucket.<bucket-name>.secret.key,Spark会自动为对应桶使用指定凭证。
方法2:读取后立即缓存数据
在读取数据后调用cache()或persist(),并触发Action(如count())将数据加载到内存/磁盘,后续写入操作直接使用缓存的数据,不再访问原S3桶。
方法3:使用独立SparkSession(不推荐)
为读取和写入分别创建独立的SparkSession,但会额外消耗资源,不适合大数据场景。
推荐方案代码实现
from pyspark.sql import SparkSession def create_spark_session(): spark = (SparkSession.builder .config("spark.hadoop.fs.s3a.fast.upload", True) .config("spark.hadoop.fs.s3a.impl", "org.apache.hadoop.fs.s3a.S3AFileSystem") .config("spark.sql.autoBroadcastJoinThreshold", -1) .config("spark.sql.shuffle.partitions", "1000") .config("spark.sql.adaptive.enabled", "true") .config("spark.sql.adaptive.coalescePartitions.enabled", "true") .config("spark.sql.adaptive.advisoryPartitionSizeInBytes", "268435456") .config("spark.jars.packages", "org.apache.hadoop:hadoop-aws:3.2.0") # 为两个桶分别配置独立凭证 .config("spark.hadoop.fs.s3a.bucket.bucket1.access.key", 'credential_bucket_1') .config("spark.hadoop.fs.s3a.bucket.bucket1.secret.key", 'credential_bucket_1') .config("spark.hadoop.fs.s3a.bucket.bucket2.access.key", 'credential_bucket_2') .config("spark.hadoop.fs.s3a.bucket.bucket2.secret.key", 'credential_bucket_2') .enableHiveSupport().getOrCreate() ) spark.sparkContext.setLogLevel("WARN") return spark def read_from_bucket1(spark): print('Running reading bucket 1') # 无需修改全局配置,直接读取 spark.read.parquet('s3a://bucket1/path/2022/*/*/').registerTempTable('temp_table') def write_to_bucket2(spark): print('Running writing bucket 2') # 无需修改全局配置,直接写入 (spark.sql('select cast(col as date) from temp_table') .write.option("compression","GZIP") .mode("overwrite") .save(f"s3a://bucket2/path_in_other_aws/")) spark = create_spark_session() read_from_bucket1(spark) write_to_bucket2(spark)
缓存方案代码实现
from pyspark.sql import SparkSession def create_spark_session(): spark = (SparkSession.builder .config("spark.hadoop.fs.s3a.fast.upload", True) .config("spark.hadoop.fs.s3a.impl", "org.apache.hadoop.fs.s3a.S3AFileSystem") .config("spark.sql.autoBroadcastJoinThreshold", -1) .config("spark.sql.shuffle.partitions", "1000") .config("spark.sql.adaptive.enabled", "true") .config("spark.sql.adaptive.coalescePartitions.enabled", "true") .config("spark.sql.adaptive.advisoryPartitionSizeInBytes", "268435456") .config("spark.jars.packages", "org.apache.hadoop:hadoop-aws:3.2.0") .enableHiveSupport().getOrCreate() ) spark.sparkContext.setLogLevel("WARN") return spark def read_from_bucket1(spark): print('Running reading bucket 1') # 设置原桶凭证 spark._jsc.hadoopConfiguration().set("fs.s3a.access.key", 'credential_bucket_1') spark._jsc.hadoopConfiguration().set("fs.s3a.secret.key", 'credential_bucket_1') # 读取后缓存并触发Action,立即加载数据 df = spark.read.parquet('s3a://bucket1/path/2022/*/*/') df.cache() df.count() # 强制读取数据到缓存 df.registerTempTable('temp_table') def write_to_bucket2(spark): print('Running writing bucket 2') # 设置目标桶凭证 spark._jsc.hadoopConfiguration().set("fs.s3a.access.key", 'credential_bucket_2') spark._jsc.hadoopConfiguration().set("fs.s3a.secret.key", 'credential_bucket_2') # 直接使用缓存数据写入,无需再访问原桶 (spark.sql('select cast(col as date) from temp_table') .write.option("compression","GZIP") .mode("overwrite") .save(f"s3a://bucket2/path_in_other_aws/")) spark = create_spark_session() read_from_bucket1(spark) write_to_bucket2(spark)
内容的提问来源于stack exchange,提问作者Renan Nogueira
相关产品推荐
相关产品推荐

