【发布时间】:2016-06-23 10:14:26
【问题描述】:
从昨天开始,我的 pySpark 代码出现了奇怪的行为。我正在使用 PyCharm 和 Spark 1.5 在 Windows 上工作。
我在 ipython 笔记本上成功运行了以下代码(使用相同版本的 python,但在集群上)。然而,当我使用 Pycharm 在我的 Windows 环境中启动它时,我得到了这个:
from pyspark.sql import SQLContext
from pyspark import SparkContext
from pyspark import SparkConf, SparkContext
from pyspark.sql import SQLContext
# SQL / Spark context:
conf = (SparkConf().setMaster("local").setAppName("analysis"))#.set("spark.executor.memory", "2g")
sc = SparkContext(conf=conf)
sqlContext = SQLContext(sc)
# Input CSV files :
inputCsvFile = "survey.csv"
separator = ','
# read the input file into a RDD
rdd = sc.textFile(inputCsvFile).split(separator)
header = rdd.first().split(separator)
# build the Schema: (some basic functions to chreate StructType object with string as default type)
schema = dictSchemaFromColumnsList(header)
schemaDf = dictSchemaToDFSchema(schema)
# create Dataframe:
df = sqlContext.createDataFrame(rdd, schemaDf)
pprint(rdd.first())
print('\ndf.count()=' + str(df.count()))
# display
df.show()
16/06/23 11:46:32 错误执行器:阶段 1.0 中任务 0.0 中的异常 (TID 1)
java.net.SocketException:对等方重置连接:套接字写入错误 在 java.net.SocketOutputStream.socketWrite0(Native Method) 在 java.net.SocketOutputStream.socketWrite(SocketOutputStream.java:109) 在 java.net.SocketOutputStream.write(SocketOutputStream.java:153) 在 java.io.BufferedOutputStream.flushBuffer(BufferedOutputStream.java:82) 在 java.io.BufferedOutputStream.write(BufferedOutputStream.java:126) 在 java.io.DataOutputStream.write(DataOutputStream.java:107) 在 java.io.FilterOutputStream.write(FilterOutputStream.java:97) 在 org.apache.spark.api.python.PythonRDD$.writeUTF(PythonRDD.scala:590) 在 org.apache.spark.api.python.PythonRDD$.org$apache$spark$api$python$PythonRDD$$write$1(PythonRDD.scala:410) 在 org.apache.spark.api.python.PythonRDD$$anonfun$writeIteratorToStream$1.apply(PythonRDD.scala:420) 在 org.apache.spark.api.python.PythonRDD$$anonfun$writeIteratorToStream$1.apply(PythonRDD.scala:420) 在 scala.collection.Iterator$class.foreach(Iterator.scala:727) 在 scala.collection.AbstractIterator.foreach(Iterator.scala:1157) 在 org.apache.spark.api.python.PythonRDD$.writeIteratorToStream(PythonRDD.scala:420) 在 org.apache.spark.api.python.PythonRDD$WriterThread$$anonfun$run$3.apply(PythonRDD.scala:249) 在 org.apache.spark.util.Utils$.logUncaughtExceptions(Utils.scala:1699) 在 org.apache.spark.api.python.PythonRDD$WriterThread.run(PythonRDD.scala:208) 16/06/23 11:46:32 WARN TaskSetManager:在阶段 1.0 中丢失任务 0.0(TID 1, localhost): java.net.SocketException: Connection reset by peer: java.net.SocketOutputStream.socketWrite0 处的套接字写入错误(本机 方法)在 java.net.SocketOutputStream.socketWrite(SocketOutputStream.java:109) 在 java.net.SocketOutputStream.write(SocketOutputStream.java:153) 在 java.io.BufferedOutputStream.flushBuffer(BufferedOutputStream.java:82) 在 java.io.BufferedOutputStream.write(BufferedOutputStream.java:126) 在 java.io.DataOutputStream.write(DataOutputStream.java:107) 在 java.io.FilterOutputStream.write(FilterOutputStream.java:97) 在 org.apache.spark.api.python.PythonRDD$.writeUTF(PythonRDD.scala:590) 在 org.apache.spark.api.python.PythonRDD$.org$apache$spark$api$python$PythonRDD$$write$1(PythonRDD.scala:410) 在 org.apache.spark.api.python.PythonRDD$$anonfun$writeIteratorToStream$1.apply(PythonRDD.scala:420) 在 org.apache.spark.api.python.PythonRDD$$anonfun$writeIteratorToStream$1.apply(PythonRDD.scala:420) 在 scala.collection.Iterator$class.foreach(Iterator.scala:727) 在 scala.collection.AbstractIterator.foreach(Iterator.scala:1157) 在 org.apache.spark.api.python.PythonRDD$.writeIteratorToStream(PythonRDD.scala:420) 在 org.apache.spark.api.python.PythonRDD$WriterThread$$anonfun$run$3.apply(PythonRDD.scala:249) 在 org.apache.spark.util.Utils$.logUncaughtExceptions(Utils.scala:1699) 在 org.apache.spark.api.python.PythonRDD$WriterThread.run(PythonRDD.scala:208)
16/06/23 11:46:32 错误 TaskSetManager: 阶段 1.0 中的任务 0 失败 1 次;中止作业 16/06/23 11:46:32 信息 TaskSchedulerImpl:已删除 TaskSet 1.0,其任务已全部完成,来自池 16/06/23 11:46:32 INFO TaskSchedulerImpl:取消阶段 1 16/06/23 11:46:32 信息 DAGScheduler:ResultStage 1(在 PythonRDD.scala:361 上运行作业) 在 0.792 秒内失败 16/06/23 11:46:32 信息 DAGScheduler:作业 1 失败: 在 PythonRDD.scala:361 上运行作业,耗时 0.802922 秒 Traceback(最近 最后调用):文件 "C:/Users/home/PycharmProjects/pySpark_analysis/Survey_2011-2016_Analysis.py", 第 38 行,在 df = sqlContext.createDataFrame(rdd, schemaDf) 文件 "C:\Spark\spark-1.5.0-bin-hadoop2.6\python\pyspark\sql\context.py", 第 404 行,在 createDataFrame rdd, schema = self._createFromRDD(data, schema, samplingRatio) 文件 "C:\Spark\spark-1.5.0-bin-hadoop2.6\python\pyspark\sql\context.py", 第 296 行,在 _createFromRDD 中 rows = rdd.take(10) 文件“C:\Spark\spark-1.5.0-bin-hadoop2.6\python\pyspark\rdd.py”,第 1299 行, 录取 res = self.context.runJob(self, takeUpToNumLeft, p) 文件“C:\Spark\spark-1.5.0-bin-hadoop2.6\python\pyspark\context.py”,行 916,在 runJob 端口 = self._jvm.PythonRDD.runJob(self._jsc.sc(), mappedRDD._jrdd, partitions) 文件 "C:\Spark\spark-1.5.0-bin-hadoop2.6\python\lib\py4j-0.8.2.1-src.zip\py4j\java_gateway.py", 第 538 行,在 call 文件中 “C:\Spark\spark-1.5.0-bin-hadoop2.6\python\pyspark\sql\utils.py”,行 36、装饰 返回 f(*a, **kw) 文件 "C:\Spark\spark-1.5.0-bin-hadoop2.6\python\lib\py4j-0.8.2.1-src.zip\py4j\protocol.py",第 300 行,在 get_return_value py4j.protocol.Py4JJavaError: 一个错误 调用 z:org.apache.spark.api.python.PythonRDD.runJob 时发生。 :org.apache.spark.SparkException:作业因阶段失败而中止: 阶段 1.0 中的任务 0 失败 1 次,最近一次失败:丢失任务 0.0 在阶段 1.0 (TID 1, localhost): java.net.SocketException: Connection 由对等方重置:套接字写入错误 java.net.SocketOutputStream.socketWrite0(本机方法)在 java.net.SocketOutputStream.socketWrite(SocketOutputStream.java:109) 在 java.net.SocketOutputStream.write(SocketOutputStream.java:153) 在 java.io.BufferedOutputStream.flushBuffer(BufferedOutputStream.java:82) 在 java.io.BufferedOutputStream.write(BufferedOutputStream.java:126) 在 java.io.DataOutputStream.write(DataOutputStream.java:107) 在 java.io.FilterOutputStream.write(FilterOutputStream.java:97) 在 org.apache.spark.api.python.PythonRDD$.writeUTF(PythonRDD.scala:590) 在 org.apache.spark.api.python.PythonRDD$.org$apache$spark$api$python$PythonRDD$$write$1(PythonRDD.scala:410) 在 org.apache.spark.api.python.PythonRDD$$anonfun$writeIteratorToStream$1.apply(PythonRDD.scala:420) 在 org.apache.spark.api.python.PythonRDD$$anonfun$writeIteratorToStream$1.apply(PythonRDD.scala:420) 在 scala.collection.Iterator$class.foreach(Iterator.scala:727) 在 scala.collection.AbstractIterator.foreach(Iterator.scala:1157) 在 org.apache.spark.api.python.PythonRDD$.writeIteratorToStream(PythonRDD.scala:420) 在 org.apache.spark.api.python.PythonRDD$WriterThread$$anonfun$run$3.apply(PythonRDD.scala:249) 在 org.apache.spark.util.Utils$.logUncaughtExceptions(Utils.scala:1699) 在 org.apache.spark.api.python.PythonRDD$WriterThread.run(PythonRDD.scala:208)
驱动程序堆栈跟踪:在 org.apache.spark.scheduler.DAGScheduler.org$apache$spark$scheduler$DAGScheduler$$failJobAndIndependentStages(DAGScheduler.scala:1280) 在 org.apache.spark.scheduler.DAGScheduler$$anonfun$abortStage$1.apply(DAGScheduler.scala:1268) 在 org.apache.spark.scheduler.DAGScheduler$$anonfun$abortStage$1.apply(DAGScheduler.scala:1267) 在 scala.collection.mutable.ResizableArray$class.foreach(ResizableArray.scala:59) 在 scala.collection.mutable.ArrayBuffer.foreach(ArrayBuffer.scala:47) 在 org.apache.spark.scheduler.DAGScheduler.abortStage(DAGScheduler.scala:1267) 在 org.apache.spark.scheduler.DAGScheduler$$anonfun$handleTaskSetFailed$1.apply(DAGScheduler.scala:697) 在 org.apache.spark.scheduler.DAGScheduler$$anonfun$handleTaskSetFailed$1.apply(DAGScheduler.scala:697) 在 scala.Option.foreach(Option.scala:236) 在 org.apache.spark.scheduler.DAGScheduler.handleTaskSetFailed(DAGScheduler.scala:697) 在 org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.doOnReceive(DAGScheduler.scala:1493) 在 org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.onReceive(DAGScheduler.scala:1455) 在 org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.onReceive(DAGScheduler.scala:1444) 在 org.apache.spark.util.EventLoop$$anon$1.run(EventLoop.scala:48) 在 org.apache.spark.scheduler.DAGScheduler.runJob(DAGScheduler.scala:567) 在 org.apache.spark.SparkContext.runJob(SparkContext.scala:1813) 在 org.apache.spark.SparkContext.runJob(SparkContext.scala:1826) 在 org.apache.spark.SparkContext.runJob(SparkContext.scala:1839) 在 org.apache.spark.api.python.PythonRDD$.runJob(PythonRDD.scala:361) 在 org.apache.spark.api.python.PythonRDD.runJob(PythonRDD.scala) 在 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:231) 在 py4j.reflection.ReflectionEngine.invoke(ReflectionEngine.java:379) 在 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:207) 在 java.lang.Thread.run(Thread.java:745) 原因: java.net.SocketException:对等方重置连接:套接字写入错误 在 java.net.SocketOutputStream.socketWrite0(Native Method) 在 java.net.SocketOutputStream.socketWrite(SocketOutputStream.java:109) 在 java.net.SocketOutputStream.write(SocketOutputStream.java:153) 在 java.io.BufferedOutputStream.flushBuffer(BufferedOutputStream.java:82) 在 java.io.BufferedOutputStream.write(BufferedOutputStream.java:126) 在 java.io.DataOutputStream.write(DataOutputStream.java:107) 在 java.io.FilterOutputStream.write(FilterOutputStream.java:97) 在 org.apache.spark.api.python.PythonRDD$.writeUTF(PythonRDD.scala:590) 在 org.apache.spark.api.python.PythonRDD$.org$apache$spark$api$python$PythonRDD$$write$1(PythonRDD.scala:410) 在 org.apache.spark.api.python.PythonRDD$$anonfun$writeIteratorToStream$1.apply(PythonRDD.scala:420) 在 org.apache.spark.api.python.PythonRDD$$anonfun$writeIteratorToStream$1.apply(PythonRDD.scala:420) 在 scala.collection.Iterator$class.foreach(Iterator.scala:727) 在 scala.collection.AbstractIterator.foreach(Iterator.scala:1157) 在 org.apache.spark.api.python.PythonRDD$.writeIteratorToStream(PythonRDD.scala:420) 在 org.apache.spark.api.python.PythonRDD$WriterThread$$anonfun$run$3.apply(PythonRDD.scala:249) 在 org.apache.spark.util.Utils$.logUncaughtExceptions(Utils.scala:1699) 在 org.apache.spark.api.python.PythonRDD$WriterThread.run(PythonRDD.scala:208)
16/06/23 11:46:32 信息 SparkContext:从关机调用 stop() 钩子
奇怪的是,如果我在调试模式下运行代码并添加如下基本指令:
People=["1,Maj,123","2,Pvt,333","3,Col,999"]
rrd1=sc.parallelize(People)
rrd1.first()
有时我的代码可以工作....这使得运行不一致.... 任何建议将不胜感激...
更新: 回顾问题后,它看起来与 Matei 之后描述的行为完全相同。显然,在缩短输入 csv 文件时问题得到了解决。
【问题讨论】:
标签: python apache-spark dataframe pyspark socketexception