【问题标题】:Accessing S3 bucket from local pyspark using assume role使用承担角色从本地 pyspark 访问 S3 存储桶
【发布时间】:2018-10-23 06:07:48
【问题描述】:

背景:为了让开发人员能够在易于使用的环境中构建和单元测试代码,我们构建了一个本地 Spark 环境,并集成了其他工具。但是,我们还想从本地环境访问 S3 和 Kinesis。当我们使用假设角色(根据我们的安全标准)从本地 Pyspark 应用程序访问 S3 时,它会引发禁止错误。

仅供参考 - 我们采用以下访问模式来访问 AWS 账户上的资源。 assume-role access pattern

测试代码access-s3-from-pyspark.py

from pyspark import SparkConf, SparkContext

conf = SparkConf().setAppName("s3a_test").setMaster("local[1]")
sc = SparkContext(conf=conf)

sc.setSystemProperty("com.amazonaws.services.s3.enableV4", "true")

hadoopConf = {}
iterator = sc._jsc.hadoopConfiguration().iterator()
while iterator.hasNext():
    prop = iterator.next()
    hadoopConf[prop.getKey()] = prop.getValue()
for item in sorted(hadoopConf.items()):
    if "fs.s3" in item[0] :
    print(item)

path="s3a://<your bucket>/test-file.txt"

## read the file for testing
lines = sc.textFile(path)

if lines.isEmpty() == False:
    lines.saveAsTextFile("test-file2.text")

属性文件spark-s3.properties

spark.hadoop.fs.s3a.impl org.apache.hadoop.fs.s3a.S3AFileSystem
spark.hadoop.fs.s3a.endpoint s3.eu-central-1.amazonaws.com
spark.hadoop.fs.s3a.access.key <your access key >
spark.hadoop.fs.s3a.secret.key <your secret key>
spark.hadoop.fs.s3a.assumed.role.sts.endpoint sts.eu-central-1.amazonaws.com
spark.hadoop.fs.s3a.aws.credentials.provider org.apache.hadoop.fs.s3a.TemporaryAWSCredentialsProvider 
spark.hadoop.fs.s3a.aws.credentials.provider org.apache.hadoop.fs.s3a.AssumedRoleCredentialProvider
spark.hadoop.fs.s3a.aws.credentials.provider org.apache.hadoop.fs.s3a.auth.AssumedRoleCredentialProvider
spark.hadoop.fs.s3a.assumed.role.session.name testSession1
spark.haeoop.fs.s3a.assumed.role.session.duration 3600
spark.hadoop.fs.s3a.assumed.role.arn <role arn>
spark.hadoop.fs.s3.canned.acl BucketOwnerFullControl

如何运行代码:

spark-submit --properties-file spark-s3.properties \
        --jars jars/hadoop-aws-2.7.3.jar,jars/aws-java-sdk-1.7.4.jar \
        access-s3-from-pyspark.pyenter code here

以上代码返回以下错误,请注意我可以通过 CLI 和 boto3 使用假设角色配置文件或 api 访问 S3。

com.amazonaws.services.s3.model.AmazonS3Exception: Status Code: 403, AWS Service: Amazon S3, AWS Request ID: 66FB4D6351898F33, AWS Error Code: null, AWS Error Message: Forbidden, S3 Extended Request ID: J8lZ4qTZ25+a8/R3ZeBTrW5TDHzo98A9iUshbe0/7VcHmiaSXZ5u6fa0TvA3E7ZYvhqXj40tf74=
    at com.amazonaws.http.AmazonHttpClient.handleErrorResponse(AmazonHttpClient.java:798)
    at com.amazonaws.http.AmazonHttpClient.executeHelper(AmazonHttpClient.java:421)
    at com.amazonaws.http.AmazonHttpClient.execute(AmazonHttpClient.java:232)
    at com.amazonaws.services.s3.AmazonS3Client.invoke(AmazonS3Client.java:3528)
    at com.amazonaws.services.s3.AmazonS3Client.getObjectMetadata(AmazonS3Client.java:976)
    at com.amazonaws.services.s3.AmazonS3Client.getObjectMetadata(AmazonS3Client.java:956)
    at org.apache.hadoop.fs.s3a.S3AFileSystem.getFileStatus(S3AFileSystem.java:892)
    at org.apache.hadoop.fs.s3a.S3AFileSystem.getFileStatus(S3AFileSystem.java:77)
    at org.apache.hadoop.fs.Globber.getFileStatus(Globber.java:57)
    at org.apache.hadoop.fs.Globber.glob(Globber.java:252)
    at org.apache.hadoop.fs.FileSystem.globStatus(FileSystem.java:1676)
    at org.apache.hadoop.mapred.FileInputFormat.singleThreadedListStatus(FileInputFormat.java:259)
    at org.apache.hadoop.mapred.FileInputFormat.listStatus(FileInputFormat.java:229)
    at org.apache.hadoop.mapred.FileInputFormat.getSplits(FileInputFormat.java:315)
    at org.apache.spark.rdd.HadoopRDD.getPartitions(HadoopRDD.scala:200)
    at org.apache.spark.rdd.RDD$$anonfun$partitions$2.apply(RDD.scala:253)
    at org.apache.spark.rdd.RDD$$anonfun$partitions$2.apply(RDD.scala:251)
    at scala.Option.getOrElse(Option.scala:121)
    at org.apache.spark.rdd.RDD.partitions(RDD.scala:251)
    at org.apache.spark.rdd.MapPartitionsRDD.getPartitions(MapPartitionsRDD.scala:35)
    at org.apache.spark.rdd.RDD$$anonfun$partitions$2.apply(RDD.scala:253)
    at org.apache.spark.rdd.RDD$$anonfun$partitions$2.apply(RDD.scala:251)
    at scala.Option.getOrElse(Option.scala:121)
    at org.apache.spark.rdd.RDD.partitions(RDD.scala:251)
    at org.apache.spark.api.java.JavaRDDLike$class.partitions(JavaRDDLike.scala:61)
    at org.apache.spark.api.java.AbstractJavaRDDLike.partitions(JavaRDDLike.scala:45)
    at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method)
    at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62)
    at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
    at java.lang.reflect.Method.invoke(Method.java:498)
    at py4j.reflection.MethodInvoker.invoke(MethodInvoker.java:244)
    at py4j.reflection.ReflectionEngine.invoke(ReflectionEngine.java:357)
    at py4j.Gateway.invoke(Gateway.java:282)
    at py4j.commands.AbstractCommand.invokeMethod(AbstractCommand.java:132)
    at py4j.commands.CallCommand.execute(CallCommand.java:79)
    at py4j.GatewayConnection.run(GatewayConnection.java:238)
    at java.lang.Thread.run(Thread.java:748)

问题:

这是正确的做法吗?

有没有其他简单的方法可以在本地使用 AWS 资源进行开发和测试(我还探索了 localstack 包,它在大多数情况下都可以工作,但仍然不完全可靠)

我是否为此使用了正确的罐子?

【问题讨论】:

  • 确保您的帐户确实有权访问对象的元数据。

标签: apache-spark amazon-s3 pyspark


【解决方案1】:

spark.hadoop.fs.s3a.aws.credentials.provider 的配置是错误的。

  • 应该只有一个条目,并且应该在一个条目中列出所有 AWS 凭据提供程序
  • S3A 假定角色提供程序(需要完全登录并要求假定角色)仅适用于最近的 Hadoop 版本 (3.1+),而不是 2.7.x,并且可能无法满足您的要求。它主要用于动态创建具有受限权限的登录并验证 S3A 连接器本身 can cope with things

您的组织对安全性要求严格,这很好,只是让生活稍微复杂了一点。

假设您可以获取帐户 ID、会话令牌和会话密钥(以某种方式),

那么对于 Hadoop 2.8+,您可以使用此填充 spark 默认值

spark.hadoop.fs.s3a.access.key AAAIKIAAA spark.hadoop.fs.s3a.secret.key ABCD spark.hadoop.fs.s3a.fs.s3a.session.token REALLYREALLYLONGVALUE spark.hadoop.fs.s3a.aws.credentials.provider org.apache.hadoop.fs.s3a.TemporaryAWSCredentialsProvider

只要会话持续,您就需要创建那些假定的角色会话机密,这曾经是 PITA,因为它们的生命是

hadoop 2.7.x 版本没有那个 TemporaryAWSCredentialsProvider,所以你必须这样做

  1. 依靠 env var 凭据提供程序,它会查找 AWS_ 环境变量。这是默认启用的,因此您根本不需要使用spark.hadoop.fs.s3a.aws.credentials.provider
  2. 将所有三个环境变量(AWS_ACCESS_KEYAWS_SECRET_KEYAWS_SESSION_TOKEN (?))设置为您从假定角色 API 调用中获得的值。
  3. 然后运行您的作业。恐怕您可能需要在任何地方设置这些环境变量。

【讨论】:

  • 感谢史蒂夫 Loughran。您的回答清楚地说明了哪个版本具有我正在寻找的功能。我已经做了以下解决:1)安装了Hadoop 2.8.4(它带有所有依赖的Jars)并配置了spark来使用它。 2) 在 core-site.xml 中配置凭据
猜你喜欢
  • 2021-12-06
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2016-02-27
  • 2015-11-02
  • 2020-06-24
相关资源
最近更新 更多