I have an error in pyspark when i am using action function after make a transformation

Viewed 32

when I use collect function after I make filter it gives me this error but if I use collect alone without any transformation function before it, running without any error

 words = sc.parallelize (
   ["scala", 
   "java", 
   "hadoop", 
   "spark", 
   "akka",
   "spark vs hadoop", 
   "pyspark",
   "pyspark and spark"]
)
words_filter = words.filter(lambda x: 'spark' in x)
filtered = words_filter.collect()
print ("Fitered RDD -> %s" % (filtered)) 

this is the error




Py4JJavaError: An error occurred while calling z:org.apache.spark.api.python.PythonRDD.collectAndServe.
: org.apache.spark.SparkException: Job aborted due to stage failure: Task 0 in stage 1.0 failed 1 times, most recent failure: Lost task 0.0 in stage 1.0 (TID 1) (DESKTOP-KKD8VJN executor driver): org.apache.spark.SparkException: Python worker failed to connect back.
   at org.apache.spark.api.python.PythonWorkerFactory.createSimpleWorker(PythonWorkerFactory.scala:189)
   at org.apache.spark.api.python.PythonWorkerFactory.create(PythonWorkerFactory.scala:109)
   at org.apache.spark.SparkEnv.createPythonWorker(SparkEnv.scala:124)
   at org.apache.spark.api.python.BasePythonRunner.compute(PythonRunner.scala:164)
   at org.apache.spark.api.python.PythonRDD.compute(PythonRDD.scala:65)
   at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:365)
   at org.apache.spark.rdd.RDD.iterator(RDD.scala:329)
   at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:90)
   at org.apache.spark.scheduler.Task.run(Task.scala:136)
   at org.apache.spark.executor.Executor$TaskRunner.$anonfun$run$3(Executor.scala:548)
   at org.apache.spark.util.Utils$.tryWithSafeFinally(Utils.scala:1504)
   at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:551)
   at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1136)
   at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:635)
   at java.base/java.lang.Thread.run(Thread.java:833)
Caused by: java.net.SocketTimeoutException: Accept timed out
   at java.base/sun.nio.ch.NioSocketImpl.timedAccept(NioSocketImpl.java:705)
   at java.base/sun.nio.ch.NioSocketImpl.accept(NioSocketImpl.java:749)
   at java.base/java.net.ServerSocket.implAccept(ServerSocket.java:673)
   at java.base/java.net.ServerSocket.platformImplAccept(ServerSocket.java:639)
   at java.base/java.net.ServerSocket.implAccept(ServerSocket.java:615)
   at java.base/java.net.ServerSocket.implAccept(ServerSocket.java:572)
   at java.base/java.net.ServerSocket.accept(ServerSocket.java:530)
   at org.apache.spark.api.python.PythonWorkerFactory.createSimpleWorker(PythonWorkerFactory.scala:176)
   ... 14 more

Driver stacktrace:
   at org.apache.spark.scheduler.DAGScheduler.failJobAndIndependentStages(DAGScheduler.scala:2672)
   at org.apache.spark.scheduler.DAGScheduler.$anonfun$abortStage$2(DAGScheduler.scala:2608)
   at org.apache.spark.scheduler.DAGScheduler.$anonfun$abortStage$2$adapted(DAGScheduler.scala:2607)
   at scala.collection.mutable.ResizableArray.foreach(ResizableArray.scala:62)
   at scala.collection.mutable.ResizableArray.foreach$(ResizableArray.scala:55)
   at scala.collection.mutable.ArrayBuffer.foreach(ArrayBuffer.scala:49)
   at org.apache.spark.scheduler.DAGScheduler.abortStage(DAGScheduler.scala:2607)
   at org.apache.spark.scheduler.DAGScheduler.$anonfun$handleTaskSetFailed$1(DAGScheduler.scala:1182)
   at org.apache.spark.scheduler.DAGScheduler.$anonfun$handleTaskSetFailed$1$adapted(DAGScheduler.scala:1182)
   at scala.Option.foreach(Option.scala:407)
   at org.apache.spark.scheduler.DAGScheduler.handleTaskSetFailed(DAGScheduler.scala:1182)
   at org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.doOnReceive(DAGScheduler.scala:2860)
   at org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.onReceive(DAGScheduler.scala:2802)
   at org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.onReceive(DAGScheduler.scala:2791)
   at org.apache.spark.util.EventLoop$$anon$1.run(EventLoop.scala:49)
   at org.apache.spark.scheduler.DAGScheduler.runJob(DAGScheduler.scala:952)
   at org.apache.spark.SparkContext.runJob(SparkContext.scala:2228)
   at org.apache.spark.SparkContext.runJob(SparkContext.scala:2249)
   at org.apache.spark.SparkContext.runJob(SparkContext.scala:2268)
   at org.apache.spark.SparkContext.runJob(SparkContext.scala:2293)
   at org.apache.spark.rdd.RDD.$anonfun$collect$1(RDD.scala:1021)
   at org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:151)
   at org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:112)
   at org.apache.spark.rdd.RDD.withScope(RDD.scala:406)
   at org.apache.spark.rdd.RDD.collect(RDD.scala:1020)
   at org.apache.spark.api.python.PythonRDD$.collectAndServe(PythonRDD.scala:180)
   at org.apache.spark.api.python.PythonRDD.collectAndServe(PythonRDD.scala)
   at java.base/jdk.internal.reflect.DirectMethodHandleAccessor.invoke(DirectMethodHandleAccessor.java:104)
   at java.base/java.lang.reflect.Method.invoke(Method.java:577)
   at py4j.reflection.MethodInvoker.invoke(MethodInvoker.java:244)
   at py4j.reflection.ReflectionEngine.invoke(ReflectionEngine.java:357)
   at py4j.Gateway.invoke(Gateway.java:282)
   at py4j.commands.AbstractCommand.invokeMethod(AbstractCommand.java:132)
   at py4j.commands.CallCommand.execute(CallCommand.java:79)
   at py4j.ClientServerConnection.waitForCommands(ClientServerConnection.java:182)
   at py4j.ClientServerConnection.run(ClientServerConnection.java:106)
   at java.base/java.lang.Thread.run(Thread.java:833)
Caused by: org.apache.spark.SparkException: Python worker failed to connect back.
   at org.apache.spark.api.python.PythonWorkerFactory.createSimpleWorker(PythonWorkerFactory.scala:189)
   at org.apache.spark.api.python.PythonWorkerFactory.create(PythonWorkerFactory.scala:109)
   at org.apache.spark.SparkEnv.createPythonWorker(SparkEnv.scala:124)
   at org.apache.spark.api.python.BasePythonRunner.compute(PythonRunner.scala:164)
   at org.apache.spark.api.python.PythonRDD.compute(PythonRDD.scala:65)
   at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:365)
   at org.apache.spark.rdd.RDD.iterator(RDD.scala:329)
   at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:90)
   at org.apache.spark.scheduler.Task.run(Task.scala:136)
   at org.apache.spark.executor.Executor$TaskRunner.$anonfun$run$3(Executor.scala:548)
   at org.apache.spark.util.Utils$.tryWithSafeFinally(Utils.scala:1504)
   at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:551)
   at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1136)
   at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:635)
   ... 1 more
Caused by: java.net.SocketTimeoutException: Accept timed out
   at java.base/sun.nio.ch.NioSocketImpl.timedAccept(NioSocketImpl.java:705)
   at java.base/sun.nio.ch.NioSocketImpl.accept(NioSocketImpl.java:749)
   at java.base/java.net.ServerSocket.implAccept(ServerSocket.java:673)
   at java.base/java.net.ServerSocket.platformImplAccept(ServerSocket.java:639)
   at java.base/java.net.ServerSocket.implAccept(ServerSocket.java:615)
   at java.base/java.net.ServerSocket.implAccept(ServerSocket.java:572)
   at java.base/java.net.ServerSocket.accept(ServerSocket.java:530)
   at org.apache.spark.api.python.PythonWorkerFactory.createSimpleWorker(PythonWorkerFactory.scala:176)
   ... 14 more

it's only appear when i am using action function after transformation function so i need you help to solve this error

0 Answers
Related