【发布时间】:2016-12-30 23:00:12
【问题描述】:
我刚开始学习 pyspark,这似乎是个大问题:我尝试将本地文本文件加载到 spark 中:
base_df = sqlContext.read.text("/root/Downloads/SogouQ1.txt")
16/12/29 11:55:20 INFO text.TextRelation:在驱动程序上列出 hdfs://localhost:9000/root/Downloads/SogouQ1.txt
base_df.show(10)
16/12/29 11:55:36 INFO storage.MemoryStore: 块 broadcast_2 已存储 作为内存中的值(估计大小 61.8 KB,空闲 78.0 KB)16/12/29 11:55:36 INFO storage.MemoryStore: 块 broadcast_2_piece0 存储为 内存中的字节数(估计大小 19.6 KB,空闲 97.6 KB)16/12/29 11:55:36 INFO storage.BlockManagerInfo:添加了 broadcast_2_piece0 本地主机上的内存:35556(大小:19.6 KB,免费:511.1 MB)16/12/29 11:55:36 INFO spark.SparkContext:从 showString 创建广播 2 在 NativeMethodAccessorImpl.java:-2 16/12/29 11:55:36 信息 storage.MemoryStore:块 broadcast_3 存储为内存中的值 (估计大小 212.1 KB,免费 309.7 KB) 16/12/29 11:55:36 INFO storage.MemoryStore: 块broadcast_3_piece0存储为字节 内存(估计大小 19.6 KB,空闲 329.2 KB)16/12/29 11:55:36 INFO storage.BlockManagerInfo:在内存中添加了 broadcast_3_piece0 本地主机:35556(大小:19.6 KB,免费:511.1 MB)16/12/29 11:55:36 信息 spark.SparkContext:从 showString 创建广播 3 NativeMethodAccessorImpl.java:-2 Traceback(最近一次调用最后一次):
文件“”,第 1 行,在文件中 “/opt/spark/python/pyspark/sql/dataframe.py”,第 257 行,显示 打印(self._jdf.showString(n,截断))文件“/opt/spark/python/lib/py4j-0.9-src.zip/py4j/java_gateway.py”,行 813,在调用文件“/opt/spark/python/pyspark/sql/utils.py”中,行 45,在装饰 返回 f(*a, **kw) 文件“/opt/spark/python/lib/py4j-0.9-src.zip/py4j/protocol.py”,第 308 行, 在 get_return_value py4j.protocol.Py4JJavaError: 发生错误 在调用 o34.showString 时。 : java.io.IOException: 没有输入路径 在工作中指定 org.apache.hadoop.mapred.FileInputFormat.listStatus(FileInputFormat.java:201) 在 org.apache.hadoop.mapred.FileInputFormat.getSplits(FileInputFormat.java:313) 在 org.apache.spark.rdd.HadoopRDD.getPartitions(HadoopRDD.scala:199) 在 org.apache.spark.rdd.RDD$$anonfun$partitions$2.apply(RDD.scala:239) 在 org.apache.spark.rdd.RDD$$anonfun$partitions$2.apply(RDD.scala:237) 在 scala.Option.getOrElse(Option.scala:120) 在 org.apache.spark.rdd.RDD.partitions(RDD.scala:237) 在 org.apache.spark.rdd.MapPartitionsRDD.getPartitions(MapPartitionsRDD.scala:35) 在 org.apache.spark.rdd.RDD$$anonfun$partitions$2.apply(RDD.scala:239) 在 org.apache.spark.rdd.RDD$$anonfun$partitions$2.apply(RDD.scala:237) 在 scala.Option.getOrElse(Option.scala:120) 在 org.apache.spark.rdd.RDD.partitions(RDD.scala:237) 在 org.apache.spark.rdd.MapPartitionsRDD.getPartitions(MapPartitionsRDD.scala:35) 在 org.apache.spark.rdd.RDD$$anonfun$partitions$2.apply(RDD.scala:239) 在 org.apache.spark.rdd.RDD$$anonfun$partitions$2.apply(RDD.scala:237) 在 scala.Option.getOrElse(Option.scala:120) 在 org.apache.spark.rdd.RDD.partitions(RDD.scala:237) 在 org.apache.spark.rdd.MapPartitionsRDD.getPartitions(MapPartitionsRDD.scala:35) 在 org.apache.spark.rdd.RDD$$anonfun$partitions$2.apply(RDD.scala:239) 在 org.apache.spark.rdd.RDD$$anonfun$partitions$2.apply(RDD.scala:237) 在 scala.Option.getOrElse(Option.scala:120) 在 org.apache.spark.rdd.RDD.partitions(RDD.scala:237) 在 org.apache.spark.sql.execution.SparkPlan.executeTake(SparkPlan.scala:190) 在 org.apache.spark.sql.execution.Limit.executeCollect(basicOperators.scala:165) 在 org.apache.spark.sql.execution.SparkPlan.executeCollectPublic(SparkPlan.scala:174) 在 org.apache.spark.sql.DataFrame$$anonfun$org$apache$spark$sql$DataFrame$$execute$1$1.apply(DataFrame.scala:1499) 在 org.apache.spark.sql.DataFrame$$anonfun$org$apache$spark$sql$DataFrame$$execute$1$1.apply(DataFrame.scala:1499) 在 org.apache.spark.sql.execution.SQLExecution$.withNewExecutionId(SQLExecution.scala:56) 在 org.apache.spark.sql.DataFrame.withNewExecutionId(DataFrame.scala:2086) 在 org.apache.spark.sql.DataFrame.org$apache$spark$sql$DataFrame$$execute$1(DataFrame.scala:1498) 在 org.apache.spark.sql.DataFrame.org$apache$spark$sql$DataFrame$$collect(DataFrame.scala:1505) 在 org.apache.spark.sql.DataFrame$$anonfun$head$1.apply(DataFrame.scala:1375) 在 org.apache.spark.sql.DataFrame$$anonfun$head$1.apply(DataFrame.scala:1374) 在 org.apache.spark.sql.DataFrame.withCallback(DataFrame.scala:2099) 在 org.apache.spark.sql.DataFrame.head(DataFrame.scala:1374) 在 org.apache.spark.sql.DataFrame.take(DataFrame.scala:1456) 在 org.apache.spark.sql.DataFrame.showString(DataFrame.scala:170) 在 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:497) 在 py4j.reflection.MethodInvoker.invoke(MethodInvoker.java:231) 在 py4j.reflection.ReflectionEngine.invoke(ReflectionEngine.java:381) 在 py4j.Gateway.invoke(Gateway.java:259) 在 py4j.commands.AbstractCommand.invokeMethod(AbstractCommand.java:133) 在 py4j.commands.CallCommand.execute(CallCommand.java:79) 在 py4j.GatewayConnection.run(GatewayConnection.java:209) 在 java.lang.Thread.run(Thread.java:745)
对于 StackOverflow 中显示的混乱错误消息,我深表歉意,我不知道如何美化它。
当我这样做时,它会起作用:
wordsDF = sqlContext.createDataFrame([('cat',), ('elephant',), ('rat',), ('rat',), ('cat',)],['word'])
wordsDF.show()
+--------+
| word|
+--------+
| cat|
|elephant|
| rat|
| rat|
| cat|
+--------+
非常感谢。
【问题讨论】:
标签: apache-spark pyspark