【问题标题】:coding reduceByKey(lambda) in map does'nt work pySpark在地图中编码reduceByKey(lambda)不起作用pySpark
【发布时间】: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


【解决方案1】:

您收到错误提示

TypeError: +: 'int' 和 'str' 的操作数类型不受支持

因为您的元组值是字符串,即 ("1,0") 而不是 (1,0),python 目前不会应用此运算符 + 或添加 intstr(string) 数据类型。

此外,在您拥有"word1" and "word2" in x 的地图函数中进行比较时似乎存在逻辑错误,因为这只会检查"word2" 是否在x 中。我会推荐以下重写:

def lan_map(x):
    if "word1" in x 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))

或者可能更短

def lan_map(x):
     return ("Count", (
         1 if "word1" in x else 0,
         1 if "word2" in x else 0
     ))

让我知道这是否适合你。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2018-08-07
    • 1970-01-01
    • 2017-11-01
    • 1970-01-01
    • 1970-01-01
    • 2020-12-02
    • 2019-06-10
    相关资源
    最近更新 更多