【发布时间】: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 更改为 s3 或 s3n,它将要求提供 aws 访问密钥。不过,我已经在 IAM 中给了 ec2 实例AmazonS3FullAccess。
IllegalArgumentException: 'AWS 访问密钥 ID 和秘密访问密钥 必须通过设置 fs.s3.awsAccessKeyId 和 fs.s3.awsSecretAccessKey 属性(分别)。'
任何帮助将不胜感激。
【问题讨论】:
标签: python apache-spark amazon-s3 pyspark