【发布时间】:2021-11-07 17:08:48
【问题描述】:
我不明白为什么我的代码不起作用。最后一行是问题:
import findspark
findspark.init()
from pyspark import SparkConf, SparkContext
from pyspark.sql.types import StringType
from pyspark import SQLContext
conf=SparkConf().setMaster("local").setAppName("mein soft")
sc=SparkContext(conf=conf)
sqlContext=SQLContext(sc)
lines=sc.textFile("File.txt")
#lines.repartition(3)
lines.getNumPartitions()
def lan_map(x):
if "word1" and "word2" in x:
return ("Count",(1,1))
elif "word1" in x:
return ("Count",("1,0"))
elif "word2" in x:
return ("Count",("0,1"))
else:
return ("Count",("0,0"))
mapfun=lines.map(lan_map)
mapfun.reduceByKey(lambda x, y: (x[0]+y[0], x[1]+y[1])).collect()
还有错误:
----------------------------------- ---------------------------- Py4JJavaError Traceback(最近调用 最后)在 1 #Esto resume lo que se hicimos 3 celdas atrás。 ----> 2 mapfun.reduceByKey(lambda x,y: (x[0]+y[0], x[1]+y[1])).collect() 3 4 #mapfun.reduceByKey(noMeFuncaLambdaAsiQueHagoEsto(mapfun.x,mupfun.y)).collect() 5 #Esto nos devuelve directamente el recuento de cuántas veces aparece "Python" y cuántas aparece "Spark"
C:\spark-3.1.2-bin-hadoop3.2\python\pyspark\rdd.py in collect(self) 第947章 948 与 SCCallSiteSync(self.context) 作为 css: --> 949 sock_info = self.ctx._jvm.PythonRDD.collectAndServe(self._jrdd.rdd()) 950 返回列表(_load_from_socket(sock_info,self._jrdd_deserializer)) 第951章
C:\spark-3.1.2-bin-hadoop3.2\python\lib\py4j-0.10.9-src.zip\py4j\java_gateway.py 在 调用(self, *args) 1302 1303 回答 = self.gateway_client.send_command(命令) -> 1304 return_value = get_return_value(1305 答案,self.gateway_client,self.target_id,self.name)1306
C:\spark-3.1.2-bin-hadoop3.2\python\pyspark\sql\utils.py in deco(*a, **千瓦) 109 def deco(*a, **kw): 110尝试: --> 111 返回 f(*a, **kw) 112 除了 py4j.protocol.Py4JJavaError 作为 e: 113 转换 = convert_exception(e.java_exception)
C:\spark-3.1.2-bin-hadoop3.2\python\lib\py4j-0.10.9-src.zip\py4j\protocol.py 在 get_return_value(answer, gateway_client, target_id, name) 324 值 = OUTPUT_CONVERTER[类型](答案[2:],网关客户端) 325 如果答案 [1] == REFERENCE_TYPE: --> 326 引发 Py4JJavaError( 327 “调用 {0}{1}{2} 时出错。\n”。 328 格式(target_id, ".", 名称), 值)
Py4JJavaError: 调用时出错 z:org.apache.spark.api.python.PythonRDD.collectAndServe。 : org.apache.spark.SparkException:作业因阶段失败而中止: 阶段 0.0 中的任务 0 失败 1 次,最近一次失败:丢失任务 0.0 在阶段 0.0 (TID 0)(LAPTOP-PB7QDPVE 执行器驱动程序): org.apache.spark.api.python.PythonException:回溯(最近 最后调用):文件 "C:\spark-3.1.2-bin-hadoop3.2\python\lib\pyspark.zip\pyspark\worker.py", 第 604 行,在主文件中 "C:\spark-3.1.2-bin-hadoop3.2\python\lib\pyspark.zip\pyspark\worker.py", 第 594 行,处理中的文件 “C:\spark-3.1.2-bin-hadoop3.2\python\pyspark\rdd.py”,第 2916 行,在 管道函数 返回 func(split, prev_func(split, iterator)) 文件“C:\spark-3.1.2-bin-hadoop3.2\python\pyspark\rdd.py”,第 2916 行,在 管道函数 返回 func(split, prev_func(split, iterator)) 文件“C:\spark-3.1.2-bin-hadoop3.2\python\pyspark\rdd.py”,第 418 行,在 功能 返回 f(iterator) 文件“C:\spark-3.1.2-bin-hadoop3.2\python\pyspark\rdd.py”,第 2144 行,在 本地结合 merge.mergeValues(iterator) 文件 "C:\spark-3.1.2-bin-hadoop3.2\python\lib\pyspark.zip\pyspark\shuffle.py", 第 242 行,在 mergeValues 中 d[k] = comb(d[k], v) if k in d else creator(v) 文件“C:\spark-3.1.2-bin-hadoop3.2\python\pyspark\util.py”,行73,在 包装 返回 f(*args, **kwargs) File "", line 2, in TypeError: unsupported operand type(s) for +: 'int' 和 'str'
在 org.apache.spark.api.python.BasePythonRunner$ReaderIterator.handlePythonException(PythonRunner.scala:517) 在 org.apache.spark.api.python.PythonRunner$$anon$3.read(PythonRunner.scala:652) 在 org.apache.spark.api.python.PythonRunner$$anon$3.read(PythonRunner.scala:635) 在 org.apache.spark.api.python.BasePythonRunner$ReaderIterator.hasNext(PythonRunner.scala:470) 在 org.apache.spark.InterruptibleIterator.hasNext(InterruptibleIterator.scala:37) 在 scala.collection.Iterator$GroupedIterator.fill(Iterator.scala:1209) 在 scala.collection.Iterator$GroupedIterator.hasNext(Iterator.scala:1215) 在 scala.collection.Iterator$$anon$10.hasNext(Iterator.scala:458) 在 org.apache.spark.shuffle.sort.BypassMergeSortShuffleWriter.write(BypassMergeSortShuffleWriter.java:132) 在 org.apache.spark.shuffle.ShuffleWriteProcessor.write(ShuffleWriteProcessor.scala:59) 在 org.apache.spark.scheduler.ShuffleMapTask.runTask(ShuffleMapTask.scala:99) 在 org.apache.spark.scheduler.ShuffleMapTask.runTask(ShuffleMapTask.scala:52) 在 org.apache.spark.scheduler.Task.run(Task.scala:131) 在 org.apache.spark.executor.Executor$TaskRunner.$anonfun$run$3(Executor.scala:497) 在 org.apache.spark.util.Utils$.tryWithSafeFinally(Utils.scala:1439) 在 org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:500) 在 java.util.concurrent.ThreadPoolExecutor.runWorker(未知来源) 在 java.util.concurrent.ThreadPoolExecutor$Worker.run(未知来源) 在 java.lang.Thread.run(Unknown Source)
驱动程序堆栈跟踪:在 org.apache.spark.scheduler.DAGScheduler.failJobAndIndependentStages(DAGScheduler.scala:2258) 在 org.apache.spark.scheduler.DAGScheduler.$anonfun$abortStage$2(DAGScheduler.scala:2207) 在 org.apache.spark.scheduler.DAGScheduler.$anonfun$abortStage$2$adapted(DAGScheduler.scala:2206) 在 scala.collection.mutable.ResizableArray.foreach(ResizableArray.scala:62) 在 scala.collection.mutable.ResizableArray.foreach$(ResizableArray.scala:55) 在 scala.collection.mutable.ArrayBuffer.foreach(ArrayBuffer.scala:49) 在 org.apache.spark.scheduler.DAGScheduler.abortStage(DAGScheduler.scala:2206) 在 org.apache.spark.scheduler.DAGScheduler.$anonfun$handleTaskSetFailed$1(DAGScheduler.scala:1079) 在 org.apache.spark.scheduler.DAGScheduler.$anonfun$handleTaskSetFailed$1$adapted(DAGScheduler.scala:1079) 在 scala.Option.foreach(Option.scala:407) 在 org.apache.spark.scheduler.DAGScheduler.handleTaskSetFailed(DAGScheduler.scala:1079) 在 org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.doOnReceive(DAGScheduler.scala:2445) 在 org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.onReceive(DAGScheduler.scala:2387) 在 org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.onReceive(DAGScheduler.scala:2376) 在 org.apache.spark.util.EventLoop$$anon$1.run(EventLoop.scala:49) 在 org.apache.spark.scheduler.DAGScheduler.runJob(DAGScheduler.scala:868) 在 org.apache.spark.SparkContext.runJob(SparkContext.scala:2196) 在 org.apache.spark.SparkContext.runJob(SparkContext.scala:2217) 在 org.apache.spark.SparkContext.runJob(SparkContext.scala:2236) 在 org.apache.spark.SparkContext.runJob(SparkContext.scala:2261) 在 org.apache.spark.rdd.RDD.$anonfun$collect$1(RDD.scala:1030) 在 org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:151) 在 org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:112) 在 org.apache.spark.rdd.RDD.withScope(RDD.scala:414) 在 org.apache.spark.rdd.RDD.collect(RDD.scala:1029) 在 org.apache.spark.api.python.PythonRDD$.collectAndServe(PythonRDD.scala:180) 在 org.apache.spark.api.python.PythonRDD.collectAndServe(PythonRDD.scala) 在 sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method) 在 sun.reflect.NativeMethodAccessorImpl.invoke(Unknown Source) at sun.reflect.DelegatingMethodAccessorImpl.invoke(未知来源)在 java.lang.reflect.Method.invoke(未知来源)在 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:238) 在 java.lang.Thread.run(Unknown Source) 原因: org.apache.spark.api.python.PythonException:回溯(最近 最后调用):文件 "C:\spark-3.1.2-bin-hadoop3.2\python\lib\pyspark.zip\pyspark\worker.py", 第 604 行,在主文件中 "C:\spark-3.1.2-bin-hadoop3.2\python\lib\pyspark.zip\pyspark\worker.py", 第 594 行,处理中的文件 “C:\spark-3.1.2-bin-hadoop3.2\python\pyspark\rdd.py”,第 2916 行,在 管道函数 返回 func(split, prev_func(split, iterator)) 文件“C:\spark-3.1.2-bin-hadoop3.2\python\pyspark\rdd.py”,第 2916 行,在 管道函数 返回 func(split, prev_func(split, iterator)) 文件“C:\spark-3.1.2-bin-hadoop3.2\python\pyspark\rdd.py”,第 418 行,在 功能 返回 f(iterator) 文件“C:\spark-3.1.2-bin-hadoop3.2\python\pyspark\rdd.py”,第 2144 行,在 本地结合 merge.mergeValues(iterator) 文件 "C:\spark-3.1.2-bin-hadoop3.2\python\lib\pyspark.zip\pyspark\shuffle.py", 第 242 行,在 mergeValues 中 d[k] = comb(d[k], v) if k in d else creator(v) 文件“C:\spark-3.1.2-bin-hadoop3.2\python\pyspark\util.py”,行73,在 包装 返回 f(*args, **kwargs) File "", line 2, in TypeError: unsupported operand type(s) for +: 'int' 和 'str'
在 org.apache.spark.api.python.BasePythonRunner$ReaderIterator.handlePythonException(PythonRunner.scala:517) 在 org.apache.spark.api.python.PythonRunner$$anon$3.read(PythonRunner.scala:652) 在 org.apache.spark.api.python.PythonRunner$$anon$3.read(PythonRunner.scala:635) 在 org.apache.spark.api.python.BasePythonRunner$ReaderIterator.hasNext(PythonRunner.scala:470) 在 org.apache.spark.InterruptibleIterator.hasNext(InterruptibleIterator.scala:37) 在 scala.collection.Iterator$GroupedIterator.fill(Iterator.scala:1209) 在 scala.collection.Iterator$GroupedIterator.hasNext(Iterator.scala:1215) 在 scala.collection.Iterator$$anon$10.hasNext(Iterator.scala:458) 在 org.apache.spark.shuffle.sort.BypassMergeSortShuffleWriter.write(BypassMergeSortShuffleWriter.java:132) 在 org.apache.spark.shuffle.ShuffleWriteProcessor.write(ShuffleWriteProcessor.scala:59) 在 org.apache.spark.scheduler.ShuffleMapTask.runTask(ShuffleMapTask.scala:99) 在 org.apache.spark.scheduler.ShuffleMapTask.runTask(ShuffleMapTask.scala:52) 在 org.apache.spark.scheduler.Task.run(Task.scala:131) 在 org.apache.spark.executor.Executor$TaskRunner.$anonfun$run$3(Executor.scala:497) 在 org.apache.spark.util.Utils$.tryWithSafeFinally(Utils.scala:1439) 在 org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:500) 在 java.util.concurrent.ThreadPoolExecutor.runWorker(未知来源) 在 java.util.concurrent.ThreadPoolExecutor$Worker.run(未知来源) ... 1 更多
我感到很失落,以至于我什至无法从我的 funmap 中返回一个位置。我的意思是这不应该工作:
mapfun[1]
我已经尝试过使用函数。但我更失败了:
def fun2(x,y):
x[0]+y[0]
x[1]+y[1]
mapfun.reduceByKey(fun2(x,y)).collect()
【问题讨论】:
-
欢迎来到 SO!错误似乎是
unsupported operand type(s) for +: 'int' and 'str'。您能否尝试删除lan_map的返回值的第二个元素中的"?
标签: python apache-spark pyspark lambda jupyter-notebook