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

如何将Web Identity角色的AWS凭证配置到Spark Context用于Spark Submit?

将Web Identity获取的AWS角色凭证传入Spark Context的方法

针对你当前的场景——Airflow已通过Web Identity Provider拿到Assumed Role临时凭证,且凭证未存储在环境变量中,以下是几种可行的传入Spark的方案:

方案一:手动提取凭证并通过Spark配置传入

  1. 获取完整临时凭证
    利用AWS CLI从当前Airflow环境中提取包含会话令牌的完整临时凭证:

    aws sts get-session-token
    

    执行后会返回AccessKeyId、SecretAccessKey和SessionToken三个关键值。

  2. 在spark-submit中指定配置
    将上述凭证直接通过spark-submit的配置参数传入:

    spark-submit \
      --conf spark.hadoop.fs.s3a.access.key=<临时AccessKeyId> \
      --conf spark.hadoop.fs.s3a.secret.key=<临时SecretAccessKey> \
      --conf spark.hadoop.fs.s3a.session.token=<临时SessionToken> \
      your_spark_job.py
    

    也可以在Spark代码中直接设置SparkConf:

    from pyspark import SparkConf, SparkContext
    
    conf = SparkConf()
    conf.set("spark.hadoop.fs.s3a.access.key", "<临时AccessKeyId>")
    conf.set("spark.hadoop.fs.s3a.secret.key", "<临时SecretAccessKey>")
    conf.set("spark.hadoop.fs.s3a.session.token", "<临时SessionToken>")
    sc = SparkContext(conf=conf)
    

方案二:让Spark直接复用Web Identity环境

无需手动提取凭证,让Spark继承Airflow的Web Identity认证环境,自动获取临时凭证:

  1. 确认Airflow环境的关键变量
    确保Airflow运行Spark任务时,已存在以下环境变量:

    • AWS_WEB_IDENTITY_TOKEN_FILE:Web Identity令牌文件路径
    • AWS_ROLE_ARN:已扮演的角色ARN
    • AWS_ROLE_SESSION_NAME(可选):会话名称
  2. 配置Spark启用Web Identity认证
    在spark-submit中添加Hadoop客户端的Web Identity认证配置:

    spark-submit \
      --conf spark.hadoop.fs.s3a.auth.type=web-identity \
      --conf spark.hadoop.fs.s3a.web.identity.token.file=$AWS_WEB_IDENTITY_TOKEN_FILE \
      --conf spark.hadoop.fs.s3a.web.identity.role.arn=$AWS_ROLE_ARN \
      your_spark_job.py
    

    这种方式会让Spark自动通过Web Identity机制获取并刷新临时凭证,无需手动管理过期问题。

方案三:通过AWS SDK在Spark代码中自动获取凭证

在Spark代码中调用AWS SDK,直接从当前环境提取临时凭证并设置到SparkContext:

import boto3
from pyspark import SparkConf, SparkContext

# 从当前环境获取临时凭证
session = boto3.Session()
creds = session.get_credentials().get_frozen_credentials()

# 配置SparkConf
conf = SparkConf()
conf.set("spark.hadoop.fs.s3a.access.key", creds.access_key)
conf.set("spark.hadoop.fs.s3a.secret.key", creds.secret_key)
conf.set("spark.hadoop.fs.s3a.session.token", creds.token)
sc = SparkContext(conf=conf)

注意事项

  • 确保Spark的Hadoop AWS客户端版本与AWS SDK版本兼容,避免出现认证兼容性问题
  • 若操作S3以外的AWS服务(如Redshift、EMR),需对应调整服务的凭证配置参数
  • 临时凭证存在过期时间,若任务运行时长超过有效期,推荐使用方案二自动刷新凭证

内容的提问来源于stack exchange,提问作者Deepak Janyavula

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.21 15:54:14