【问题标题】:Access public available Amazon S3 file from Apache Spark从 Apache Spark 访问公共可用的 Amazon S3 文件
【发布时间】:2018-12-17 17:40:35
【问题描述】:

我有一个公开可用的 Amazon s3 资源(文本文件)并希望从 spark 访问它。这意味着 - 我没有任何亚马逊凭证 - 如果我只想下载它,它可以正常工作:

val bucket = "<my-bucket>"
val key = "<my-key>"

val client = new AmazonS3Client
val o = client.getObject(bucket, key)
val content = o.getObjectContent // <= can be read and used as input stream

但是,当我尝试从 spark 上下文访问相同的资源时

val conf = new SparkConf().setAppName("app").setMaster("local")
val sc = new SparkContext(conf)
val f = sc.textFile(s"s3a://$bucket/$key")
println(f.count())

我收到以下堆栈跟踪错误:

Exception in thread "main" com.amazonaws.AmazonClientException: Unable to load AWS credentials from any provider in the chain
    at com.amazonaws.auth.AWSCredentialsProviderChain.getCredentials(AWSCredentialsProviderChain.java:117)
    at com.amazonaws.services.s3.AmazonS3Client.invoke(AmazonS3Client.java:3521)
    at com.amazonaws.services.s3.AmazonS3Client.headBucket(AmazonS3Client.java:1031)
    at com.amazonaws.services.s3.AmazonS3Client.doesBucketExist(AmazonS3Client.java:994)
    at org.apache.hadoop.fs.s3a.S3AFileSystem.initialize(S3AFileSystem.java:297)
    at org.apache.hadoop.fs.FileSystem.createFileSystem(FileSystem.java:2653)
    at org.apache.hadoop.fs.FileSystem.access$200(FileSystem.java:92)
    at org.apache.hadoop.fs.FileSystem$Cache.getInternal(FileSystem.java:2687)
    at org.apache.hadoop.fs.FileSystem$Cache.get(FileSystem.java:2669)
    at org.apache.hadoop.fs.FileSystem.get(FileSystem.java:371)
    at org.apache.hadoop.fs.Path.getFileSystem(Path.java:295)
    at org.apache.hadoop.mapred.FileInputFormat.listStatus(FileInputFormat.java:221)
    at org.apache.hadoop.mapred.FileInputFormat.getSplits(FileInputFormat.java:270)
    at org.apache.spark.rdd.HadoopRDD.getPartitions(HadoopRDD.scala:207)
    at org.apache.spark.rdd.RDD$$anonfun$partitions$2.apply(RDD.scala:219)
    at org.apache.spark.rdd.RDD$$anonfun$partitions$2.apply(RDD.scala:217)
    at scala.Option.getOrElse(Option.scala:121)
    at org.apache.spark.rdd.RDD.partitions(RDD.scala:217)
    at org.apache.spark.rdd.MapPartitionsRDD.getPartitions(MapPartitionsRDD.scala:32)
    at org.apache.spark.rdd.RDD$$anonfun$partitions$2.apply(RDD.scala:219)
    at org.apache.spark.rdd.RDD$$anonfun$partitions$2.apply(RDD.scala:217)
    at scala.Option.getOrElse(Option.scala:121)
    at org.apache.spark.rdd.RDD.partitions(RDD.scala:217)
    at org.apache.spark.SparkContext.runJob(SparkContext.scala:1781)
    at org.apache.spark.rdd.RDD.count(RDD.scala:1099)
    at com.example.Main$.main(Main.scala:14)
    at com.example.Main.main(Main.scala)
    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:497)
    at com.intellij.rt.execution.application.AppMain.main(AppMain.java:140)

我不想提供任何 AWS 凭证——我只想匿名访问资源(目前)——如何实现这一点?我可能需要让它使用像 AnonymousAWSCredentialsProvider 这样的东西——但是如何把它放在 spark 或 hadoop 中?

附:我的 build.sbt 以防万一

scalaVersion := "2.11.7"

libraryDependencies ++= Seq(
  "org.apache.spark" %% "spark-core" % "1.4.1",
  "org.apache.hadoop" % "hadoop-aws" % "2.7.1"
)

更新:在我做了一些调查之后 - 我明白了它不起作用的原因。

首先,S3AFileSystem 使用以下凭证顺序创建 AWS 客户端:

AWSCredentialsProviderChain credentials = new AWSCredentialsProviderChain(
    new BasicAWSCredentialsProvider(accessKey, secretKey),
    new InstanceProfileCredentialsProvider(),
    new AnonymousAWSCredentialsProvider()
);

“accessKey”和“secretKey”值取自 spark conf 实例(密钥必须是“fs.s3a.access.key”和“fs.s3a.secret.key”或 org.apache.hadoop.fs。 s3a.Constants.ACCESS_KEY 和 org.apache.hadoop.fs.s3a.Constants.SECRET_KEY 常量,更方便)。

第二 - 您可能会看到 AnonymousAWSCredentialsProvider 是第三个选项(最后一个优先级) - 这可能有什么问题?查看 AnonymousAWSCredentials 的实现:

public class AnonymousAWSCredentials implements AWSCredentials {

    public String getAWSAccessKeyId() {
        return null;
    }

    public String getAWSSecretKey() {
        return null;
    }
}

它只是为访问密钥和秘密密钥返回 null。听起来很合理。但是看看 AWSCredentialsProviderChain:

AWSCredentials credentials = provider.getCredentials();

if (credentials.getAWSAccessKeyId() != null &&
    credentials.getAWSSecretKey() != null) {
    log.debug("Loading credentials from " + provider.toString());

    lastUsedProvider = provider;
    return credentials;
}

如果两个密钥都为空,它不会选择提供者 - 这意味着匿名凭据无法工作。看起来像 aws-java-sdk-1.7.4 中的错误。我尝试使用最新版本 - 但它与 hadoop-aws-2.7.1 不兼容。

还有其他想法吗?

【问题讨论】:

  • 你有没有成功,也许是更新的版本?
  • 不,我有一段时间没有尝试过这个 - 我什至忘记了,不要使用 amazon s3 做任何事情

标签: scala amazon-s3 apache-spark


【解决方案1】:

我个人从未访问过 Spark 的公共数据。您可以尝试使用虚拟凭据,或仅为此用途创建一个。直接在 SparkConf 对象上设置它们。

val sparkConf: SparkConf = ???
val accessKeyId: String = ???
val secretAccessKey: String = ???
sparkConf.set("spark.hadoop.fs.s3.awsAccessKeyId", accessKeyId)
sparkConf.set("spark.hadoop.fs.s3n.awsAccessKeyId", accessKeyId)
sparkConf.set("spark.hadoop.fs.s3.awsSecretAccessKey", secretAccessKey)
sparkConf.set("spark.hadoop.fs.s3n.awsSecretAccessKey", secretAccessKey)

作为替代方法,请阅读 DefaultAWSCredentialsProviderChain 的文档以查看在何处查找凭据。列表(顺序很重要)是:

  • 环境变量 - AWS_ACCESS_KEY_ID 和 AWS_SECRET_KEY
  • Java 系统属性 - aws.accessKeyId 和 aws.secretKey
  • 所有 AWS 开发工具包和 AWS CLI 共享的默认位置 (~/.aws/credentials) 中的凭证配置文件
  • 通过 Amazon EC2 元数据服务交付的实例配置文件凭证

【讨论】:

  • 还是有问题。我将以下值添加到您给我的键中(确切的字符串“aaa”作为虚拟凭据)。我希望在最坏的情况下会看到身份验证错误,但我看到了相同的异常“无法从链中的任何提供商加载 AWS 凭证”
  • 正确的密钥必须是“spark.hadoop.fs.s3a.access.key”和“spark.hadoop.fs.s3a.secret.key”。顺便说一句,提供虚拟值没有t 帮助 - 现在我看到 403 错误。看起来不可能使用带有 spark 的 AWS S3 的匿名凭证。根据源代码 - 凭证的顺序是不同的。AWSCredentialsProviderChain credentials = new AWSCredentialsProviderChain(new BasicAWSCredentialsProvider(accessKey, secretKey) , new InstanceProfileCredentialsProvider(), new AnonymousAWSCredentialsProvider() ); 而匿名根本不起作用。
  • 好的,抱歉,我没有看到您使用的是s3a 协议。你试过s3n 吗?
【解决方案2】:

这对我有帮助:

val session = SparkSession.builder()
  .appName("App")
  .master("local[*]") 
  .config("fs.s3a.aws.credentials.provider", "org.apache.hadoop.fs.s3a.AnonymousAWSCredentialsProvider")
  .getOrCreate()

val df = session.read.csv(filesFromS3:_*)

版本:

"org.apache.spark" %% "spark-sql" % "2.4.0",
"org.apache.hadoop" % "hadoop-aws" % "2.8.5",

文档: https://hadoop.apache.org/docs/current/hadoop-aws/tools/hadoop-aws/index.html#Authentication_properties

【讨论】:

【解决方案3】:

看来您现在可以使用 aws.credentials.provider 配置键来使用由 org.apache.hadoop.fs.s3a.AnonymousAWSCredentialsProvider 提供的匿名访问,它正确地 special case 匿名提供者。但是,您需要比 2.7 更新的 hadoop-aws,这意味着您还需要没有捆绑 hadoop 的 spark 安装。

这是我在 colab 中的做法:

!apt-get install openjdk-8-jdk-headless -qq > /dev/null
!wget -q http://apache.osuosl.org/spark/spark-2.3.1/spark-2.3.1-bin-without-hadoop.tgz
!tar xf spark-2.3.1-bin-without-hadoop.tgz
!pip install -q findspark
!pip install -q pyarrow

现在我们在边上安装hadoop,并将hadoop classpath的输出设置为SPARK_DIST_CLASSPATH,这样spark就可以看到了。

import os
!wget -q http://mirror.nbtelecom.com.br/apache/hadoop/common/hadoop-2.8.4/hadoop-2.8.4.tar.gz
!tar xf hadoop-2.8.4.tar.gz
os.environ['HADOOP_HOME']= '/content/hadoop-2.8.4'
os.environ["SPARK_DIST_CLASSPATH"] = "/content/hadoop-2.8.4/etc/hadoop:/content/hadoop-2.8.4/share/hadoop/common/lib/*:/content/hadoop-2.8.4/share/hadoop/common/*:/content/hadoop-2.8.4/share/hadoop/hdfs:/content/hadoop-2.8.4/share/hadoop/hdfs/lib/*:/content/hadoop-2.8.4/share/hadoop/hdfs/*:/content/hadoop-2.8.4/share/hadoop/yarn/lib/*:/content/hadoop-2.8.4/share/hadoop/yarn/*:/content/hadoop-2.8.4/share/hadoop/mapreduce/lib/*:/content/hadoop-2.8.4/share/hadoop/mapreduce/*:/content/hadoop-2.8.4/contrib/capacity-scheduler/*.jar"

然后我们确实喜欢https://mikestaszel.com/2018/03/07/apache-spark-on-google-colaboratory/,但添加了 s3a 和匿名阅读支持,这就是问题所在。

import os
os.environ["JAVA_HOME"] = "/usr/lib/jvm/java-8-openjdk-amd64"
os.environ["SPARK_HOME"] = "/content/spark-2.3.1-bin-without-hadoop"
os.environ['PYSPARK_SUBMIT_ARGS'] = '--packages com.amazonaws:aws-java-sdk:1.10.6,org.apache.hadoop:hadoop-aws:2.8.4 --conf spark.sql.execution.arrow.enabled=true --conf spark.hadoop.fs.s3a.aws.credentials.provider=org.apache.hadoop.fs.s3a.AnonymousAWSCredentialsProvider pyspark-shell'

最后我们可以创建会话了。

import findspark
findspark.init()
from pyspark.sql import SparkSession
spark = SparkSession.builder.master("local[*]").getOrCreate()

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2016-02-08
    • 2015-08-23
    • 2014-08-13
    • 2019-04-20
    • 2021-04-27
    相关资源
    最近更新 更多