使用Pandas to_parquet向S3传输文件时多套凭证的访问拒绝问题
问题描述
我有一个数据处理管道,完成数据处理后将Parquet文件传输至S3。为避免本地存储,我通过在to_parquet调用中使用S3 URI直接上传,代码如下:
@task(name='upload_parquet', retries=2, retry_delay_seconds=2) def upload_to_s3(df: pd.DataFrame, bucket: str, key: str): bucket = bucket if not bucket.endswith('/') else bucket[:-1] access_key, secret_key = get_s3_credentials(bucket) os.environ['AWS_ACCESS_KEY_ID'] = access_key os.environ['AWS_SECRET_ACCESS_KEY'] = secret_key.get_secret_value() df.to_parquet(f's3://{bucket}/{key}', engine='pyarrow', compression='snappy') del os.environ['AWS_ACCESS_KEY_ID'] del os.environ['AWS_SECRET_ACCESS_KEY']
该代码在第一个Bucket中可正常运行,但切换至第二个Bucket(使用另一套凭证)时,出现AccessDenied错误。猜测可能是boto3在某一层级缓存了凭证,请问有哪些简洁的推荐处理方式?
解决方案
1. 直接传入PyArrow的S3配置(优先推荐)
PyArrow的to_parquet支持通过storage_options参数直接指定S3凭证,完全绕过环境变量,从根源避免凭证缓存问题:
@task(name='upload_parquet', retries=2, retry_delay_seconds=2) def upload_to_s3(df: pd.DataFrame, bucket: str, key: str): bucket = bucket if not bucket.endswith('/') else bucket[:-1] access_key, secret_key = get_s3_credentials(bucket) storage_options = { 'access_key': access_key, 'secret_key': secret_key.get_secret_value() } df.to_parquet( f's3://{bucket}/{key}', engine='pyarrow', compression='snappy', storage_options=storage_options )
2. 清除boto3会话缓存
如果必须依赖环境变量,每次执行前主动清除boto3和botocore的缓存会话,强制重新读取最新环境变量:
import boto3 from botocore.session import Session @task(name='upload_parquet', retries=2, retry_delay_seconds=2) def upload_to_s3(df: pd.DataFrame, bucket: str, key: str): bucket = bucket if not bucket.endswith('/') else bucket[:-1] access_key, secret_key = get_s3_credentials(bucket) # 清除缓存会话 boto3.DEFAULT_SESSION = None Session().clear_cache() os.environ['AWS_ACCESS_KEY_ID'] = access_key os.environ['AWS_SECRET_ACCESS_KEY'] = secret_key.get_secret_value() df.to_parquet(f's3://{bucket}/{key}', engine='pyarrow', compression='snappy') del os.environ['AWS_ACCESS_KEY_ID'] del os.environ['AWS_SECRET_ACCESS_KEY']
3. 使用独立boto3会话绑定凭证
创建独立的boto3会话并传递给PyArrow,确保每个Bucket的凭证上下文完全隔离:
import boto3 @task(name='upload_parquet', retries=2, retry_delay_seconds=2) def upload_to_s3(df: pd.DataFrame, bucket: str, key: str): bucket = bucket if not bucket.endswith('/') else bucket[:-1] access_key, secret_key = get_s3_credentials(bucket) session = boto3.Session( aws_access_key_id=access_key, aws_secret_access_key=secret_key.get_secret_value() ) s3_client = session.client('s3') df.to_parquet( f's3://{bucket}/{key}', engine='pyarrow', compression='snappy', storage_options={'client': s3_client} )
方案优先级
优先采用方案1,逻辑简洁直接,完全避开环境变量和boto3缓存的坑,是最可靠的处理方式。
内容的提问来源于stack exchange,提问作者filpa
相关产品推荐
相关产品推荐

