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

如何在PySpark中使用两套AWS凭证跨账号读写S3存储桶

跨AWS账号S3数据迁移的凭证冲突问题

问题背景

需要从某AWS账号的S3存储桶读取多份文件,写入另一AWS账号的S3存储桶。

现状问题

在函数中切换AWS凭证后,执行写入操作时出现读取相关错误——Spark执行写入时仍使用新凭证去访问原存储桶,导致权限不足。

问题原因

  1. Spark的Hadoop配置是全局共享的,修改spark._jsc.hadoopConfiguration()的凭证会直接覆盖全局设置。
  2. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.19 14:31:02