【问题标题】:How should I load file on s3 using Spark?我应该如何使用 Spark 在 s3 上加载文件?
【发布时间】:2018-11-03 02:57:47
【问题描述】:

我通过pip install pyspark安装了spark

我正在使用以下代码从 s3 上的文件创建数据框。

from pyspark.sql import SparkSession

spark = SparkSession.builder \
            .config('spark.driver.extraClassPath', '/home/ubuntu/spark/jars/aws-java-sdk-1.11.335.jar:/home/ubuntu/spark/jars/hadoop-aws-2.8.4.jar') \
            .appName("cluster").getOrCreate()
df = spark.read.load('s3a://bucket/path/to/file')

但是我得到了一个错误:

----------------------------------- ---------------------------- Py4JJavaError Traceback(最近调用 最后)在() ----> 1 df = spark.read.load('s3a://bucket/path/to/file')

~/miniconda3/envs/audience/lib/python3.6/site-packages/pyspark/sql/readwriter.py 在加载(自我,路径,格式,模式,**选项) 第164章 165 如果是实例(路径,基本字符串): --> 166 返回 self._df(self._jreader.load(path)) 167 elif 路径不是无: 168 如果类型(路径)!= 列表:

~/miniconda3/envs/audience/lib/python3.6/site-packages/py4j/java_gateway.py 在 调用(self, *args) 1158 答案 = self.gateway_client.send_command(command) 1159 return_value = get_return_value( -> 1160 answer, self.gateway_client, self.target_id, self.name) 1161 1162 for temp_args in temp_args:

~/miniconda3/envs/audience/lib/python3.6/site-packages/pyspark/sql/utils.py 装饰中(*a,**kw) 61 def deco(*a, **kw): 62 尝试: ---> 63 返回 f(*a, **kw) 64 除了 py4j.protocol.Py4JJavaError 作为 e: 65 秒 = e.java_exception.toString()

~/miniconda3/envs/audience/lib/python3.6/site-packages/py4j/protocol.py 在 get_return_value(answer, gateway_client, target_id, name) 第318章 319 “调用 {0}{1}{2} 时出错。\n”。 --> 320 格式(target_id, ".", name), value) 321 其他: 第322章

Py4JJavaError:调用 o220.load 时出错。 : java.lang.NoClassDefFoundError: org/apache/hadoop/fs/StorageStatistics 在 java.lang.Class.forName0(Native Method) 在 java.lang.Class.forName(Class.java:348) 在 org.apache.hadoop.conf.Configuration.getClassByNameOrNull(Configuration.java:2134) 在 org.apache.hadoop.conf.Configuration.getClassByName(Configuration.java:2099) 在 org.apache.hadoop.conf.Configuration.getClass(Configuration.java:2193) 在 org.apache.hadoop.fs.FileSystem.getFileSystemClass(FileSystem.java:2654) 在 org.apache.hadoop.fs.FileSystem.createFileSystem(FileSystem.java:2667) 在 org.apache.hadoop.fs.FileSystem.access$200(FileSystem.java:94) 在 org.apache.hadoop.fs.FileSystem$Cache.getInternal(FileSystem.java:2703) 在 org.apache.hadoop.fs.FileSystem$Cache.get(FileSystem.java:2685) 在 org.apache.hadoop.fs.FileSystem.get(FileSystem.java:373) 在 org.apache.hadoop.fs.Path.getFileSystem(Path.java:295) 在 org.apache.spark.sql.execution.streaming.FileStreamSink$.hasMetadata(FileStreamSink.scala:44) 在 org.apache.spark.sql.execution.datasources.DataSource.resolveRelation(DataSource.scala:354) 在 org.apache.spark.sql.DataFrameReader.loadV1Source(DataFrameReader.scala:239) 在 org.apache.spark.sql.DataFrameReader.load(DataFrameReader.scala:227) 在 org.apache.spark.sql.DataFrameReader.load(DataFrameReader.scala:174) 在 sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method) 在 sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62) 在 sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43) 在 java.lang.reflect.Method.invoke(Method.java:498) 在 py4j.reflection.MethodInvoker.invoke(MethodInvoker.java:244) 在 py4j.reflection.ReflectionEngine.invoke(ReflectionEngine.java:357) 在 py4j.Gateway.invoke(Gateway.java:282) 在 py4j.commands.AbstractCommand.invokeMethod(AbstractCommand.java:132) 在 py4j.commands.CallCommand.execute(CallCommand.java:79) 在 py4j.GatewayConnection.run(GatewayConnection.java:214) 在 java.lang.Thread.run(Thread.java:748) 原因: java.lang.ClassNotFoundException: org.apache.hadoop.fs.StorageStatistics 在 java.net.URLClassLoader.findClass(URLClassLoader.java:381) 在 java.lang.ClassLoader.loadClass(ClassLoader.java:424) 在 sun.misc.Launcher$AppClassLoader.loadClass(Launcher.java:349) 在 java.lang.ClassLoader.loadClass(ClassLoader.java:357) ... 28 更多

如果我将 s3a 更改为 s3s3n,它将要求提供 aws 访问密钥。不过,我已经在 IAM 中给了 ec2 实例AmazonS3FullAccess

IllegalArgumentException: 'AWS 访问密钥 ID 和秘密访问密钥 必须通过设置 fs.s3.awsAccessKeyId 和 fs.s3.awsSecretAccessKey 属性(分别)。'

任何帮助将不胜感激。

【问题讨论】:

    标签: python apache-spark amazon-s3 pyspark


    【解决方案1】:

    您需要一种将 AWS 凭证公开给脚本的方法。

    下面使用 botocore 的示例可能会超出范围,但可以让您无需滚动自己的 AWS 配置或凭证解析器。

    首先,

    pip install botocore

    然后创建一个会话并盲目地解析您的凭据。凭证解析顺序为documented here

    from pyspark.sql import SparkSession
    import botocore.session
    
    session = botocore.session.get_session()
    credentials = session.get_credentials()
    
    spark = (
        SparkSession
        .builder
        .config(
            'spark.driver.extraClassPath', 
            '/home/ubuntu/spark/jars/aws-java-sdk-1.11.335.jar:'
            '/home/ubuntu/spark/jars/hadoop-aws-2.8.4.jar')
        .config('fs.s3a.access.key', credentials.access_key)
        .config('fs.s3a.secret.key', credentials.secret_key)
        .appName("cluster")
        .getOrCreate()
    )
    
    df = spark.read.load('s3a://bucket/path/to/file')
    

    编辑

    使用s3n文件系统客户端时,authentication properties是这样的

    .config('fs.s3n.awsAccessKeyId', credentials.access_key)
    .config('fs.s3n.awsSecretAccessKey', credentials.secret_key)
    

    【讨论】:

    • 它没有用。我仍然收到相同的错误“调用 o570.load 时发生错误。:java.lang.NoClassDefFoundError:org/apache/hadoop/fs/StorageStatistics”。如果我将s3a 更改为s3n,我会得到'org.apache.hadoop.security.AccessControlException: Permission denied'。我确定我的 IAM 是正确的,因为我可以通过 boto3 访问 s3 存储桶
    【解决方案2】:

    第一个错误告诉您 Spark 尝试加载类 org.apache.hadoop.fs.StorageStatistics。你能确保你的 Spark 版本适合你的 Hadoop JAR 吗?通常,在此提交 https://github.com/apache/hadoop/commit/687233f20d24c29041929dd0a99d963cec54b6df#diff-114b1833bd381e88382ade201dc692e8 中添加了 Spark 尝试加载的类,并且关于发布标签,首先在 3.0.0 中发布。由于您使用的是 Hadoop 2.8.4,因此将其升级到 3.0.0 可能是一个解决方案。

    【讨论】:

    • 我不确定我是否使用 Hadoop,因为我通过 pipy 安装了 Spark。我想我可能正在使用独立的 Spark。火花版本是 2.3。我尝试了hadoop-aws-3.1.0,这给了我同样的错误。
    • hadoop-aws 版本与 spark 类路径上的 hadoop-common JAR 匹配 100% 至关重要,否则有效保证堆栈跟踪。不,不要尝试 hadoop 3.x 和 spark,还没有
    • 嗨@SteveLoughran,我想我根本没有使用hadoop。我的 spark 是通过 pip install pyspark 安装的,我试过 hadoop-aws 2.732.843.1,但都没有工作。我总是遇到同样的错误'java.lang.NoClassDefFoundError: org/apache/hadoop/fs/StorageStatistics'
    • Spark 使用 hadoop s3 连接器与 S3 通信。因此,您正在使用 hadoop。你有类路径问题,你必须自己识别和修复。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2017-06-15
    • 2019-10-31
    • 1970-01-01
    • 2020-05-02
    • 2017-05-29
    • 2018-09-30
    • 2010-09-05
    相关资源
    最近更新 更多