【问题标题】:Pyspark collectAsMap() UDAF alternative - Could not serialize object: Py4JError: An error occurred while calling o62.__getstate__ [duplicate]Pyspark collectAsMap()UDAF替代方案-无法序列化对象:Py4JError:调用o62.__getstate__时发生错误[重复]
【发布时间】:2019-08-24 02:09:11
【问题描述】:

我正在尝试将聚合函数应用于 pyspark 中的列。背景是我手头只有 Spark 2.2,没有使用矢量化 pandas_udf 的选项

sdt = spark.createDataFrame(zip([random.randint(1,100) for x in range(20)], [random.randint(1,100) for x in range(20)]), schema=['col1', 'col2'])
+----+----+
|col1|col2|
+----+----+
|  19|  51|
|  95|  56|
|  11|  94|
|  80|  99|
|  20|  80|
|  38|  91|
|  18|  88|
|   4|  33|
+----+----+

为了使列并行化,我将其转换为 rdd

sdt_col_rdd = sc.parallelize(sdt.columns)

使用普通 python 函数进行测试,返回 panda 的数据帧

x = sdt_col_rdd.map(lambda col : (col, pd.DataFrame(np.random.randint(0,100,size=(2, 4)), columns=list('ABCD'))))
y = x.collectAsMap() #collect into dictionary with column names as key
print(y['col1']; print(y['col2']);
    A   B   C   D
0  14  55   4  57
1  36  84  53  51
    A   B   C   D
0  14  55   4  57
1  36  84  53  51

切换到 spark 数据帧,这里也是一个返回 panda 的 df 的示例函数,但是处理 Spark 的 df 并使用它的本机聚合、转换、动作等:

def myFunc(df, c):
    #other more processing, aggregation, transformation may be performed here
    res = df.agg((F.min(c)-1).alias("min_"+c), (F.max(c)+1).alias("max_"+c)).toPandas()
    res["col_name"] = c
    return res

函数本身就可以正常工作

myFunc(sdt.select('col1'), 'col1')
    min_col1    max_col1    col_name
0   4   100 col1

当我将它放入 rdd 映射时出现问题,类似于上面所做的

x= sdt_col_rdd.map(lambda col: (col,myFunc(sdt.select(col), col)))
y = x.collectAsMap()

知道如何在 Spark 中为列并行实现这种转换/动作,而无需 udaf 吗?由于数据集庞大且不利用 Spark 的功能,Collect_list 效率不高。

在处理上述异常的过程中,又发生了一个异常:

PicklingError Traceback(最近调用 最后)在() 1 col_map = sdt_col_rdd.map(lambda col: (col,myFunc(sdt.select(col), col))) ----> 2 y = col_map.collectAsMap()

/data/2/parcels/SPARK2-2.2.0.cloudera4-1.cdh5.13.3.p0.603055/lib/spark2/python/pyspark/rdd.py 在 collectAsMap(self) 1555 4 1556 """ -> 1557 return dict(self.collect()) 1558 1559 def keys(self):

/data/2/parcels/SPARK2-2.2.0.cloudera4-1.cdh5.13.3.p0.603055/lib/spark2/python/pyspark/rdd.py in collect(self) 第794章 795 使用 SCCallSiteSync(self.context) 作为 css: --> 796 sock_info = self.ctx._jvm.PythonRDD.collectAndServe(self._jrdd.rdd()) 第797章 第798章

/data/2/parcels/SPARK2-2.2.0.cloudera4-1.cdh5.13.3.p0.603055/lib/spark2/python/pyspark/rdd.py in _jrdd(self) 2440 2441 Wrapped_func = _wrap_function(self.ctx,self.func,self._prev_jrdd_deserializer, -> 2442 self._jrdd_deserializer,分析器)2443 python_rdd = self.ctx._jvm.PythonRDD(self._prev_jrdd.rdd(), Wrapped_func, 2444
self.preservesPartitioning)

/data/2/parcels/SPARK2-2.2.0.cloudera4-1.cdh5.13.3.p0.603055/lib/spark2/python/pyspark/rdd.py in _wrap_function(sc, func, deserializer, serializer,分析器)
2373 断言序列化程序,“序列化程序不应为空”2374
command = (func, profiler, deserializer, serializer) -> 2375 pickled_command,broadcast_vars,env,包括=_prepare_for_python_RDD(sc,command)2376 return sc._jvm.PythonFunction(bytearray(pickled_command),env,包括, sc.pythonExec, 2377 sc.pythonVer, broadcast_vars, sc._javaAccumulator)

/data/2/parcels/SPARK2-2.2.0.cloudera4-1.cdh5.13.3.p0.603055/lib/spark2/python/pyspark/rdd.py in _prepare_for_python_RDD(sc, command) 2359 # 序列化 命令将被广播压缩 2360 ser = CloudPickleSerializer() -> 2361 pickled_command = ser.dumps(command) 2362 if len(pickled_command) > (1

/data/2/parcels/SPARK2-2.2.0.cloudera4-1.cdh5.13.3.p0.603055/lib/spark2/python/pyspark/serializers.py 在转储(自我,obj) 462 463 def 转储(自我,obj): --> 464 返回 cloudpickle.dumps(obj, 2) 465 第466章

/data/2/parcels/SPARK2-2.2.0.cloudera4-1.cdh5.13.3.p0.603055/lib/spark2/python/pyspark/cloudpickle.py 在转储中(obj,协议) 702 703 cp = CloudPickler(文件、协议) --> 704 cp.dump(obj) 705 706 返回文件.getvalue()

/data/2/parcels/SPARK2-2.2.0.cloudera4-1.cdh5.13.3.p0.603055/lib/spark2/python/pyspark/cloudpickle.py 在转储(自我,obj) 160 msg = "无法序列化对象:%s: %s" % (e.class.name, emsg) 第161章 --> 162 引发 pickle.PicklingError(msg) 163 164 def save_memoryview(自我,obj):

PicklingError:无法序列化对象:Py4JError:错误 调用 o62.getstate 时发生。跟踪:py4j.Py4JException: 方法 getstate([]) 不存在于 py4j.reflection.ReflectionEngine.getMethod(ReflectionEngine.java:318) 在 py4j.reflection.ReflectionEngine.getMethod(ReflectionEngine.java:326) 在 py4j.Gateway.invoke(Gateway.java:274) 在 py4j.commands.AbstractCommand.invokeMethod(AbstractCommand.java:132) 在 py4j.commands.CallCommand.execute(CallCommand.java:79) 在 py4j.GatewayConnection.run(GatewayConnection.java:238) 在 java.lang.Thread.run(Thread.java:748)

【问题讨论】:

  • 正如我所提到的,基于行的 UDF 在这种情况下不起作用。我们需要 UDAF,它在 Spark 2.2 中不可用,我正在寻找替代方案,不使用 collect_list 因为 1/ 这将失去 Spark 在该单列上的分布式能力;我们正在调整到普通自己的 Python 列表 2/collect_list 将有数百万条记录的问题
  • 似乎您没有抓住已接受答案的这一部分 - “您正在将 pyspark 数据帧、df_whitelist 传递给 UDF,pyspark 数据帧不能被腌制。您还在内部的数据帧上进行计算不可接受(不可能)的 UDF。"

标签: scala apache-spark pyspark mapreduce rdd


【解决方案1】:

好像你没有注册你的udf,导入udf函数并注册udf,如下图,应该可以的。

从 pyspark.sql.functions 导入 *

myFunc=udf(myFunc,StringType())

【讨论】:

  • 据我了解,udf 用于 Spark 数据帧,基于行。我作为 rdd 访问,因为它是基于列的,而不是数据框。此外,我认为 udf 不适用于聚合函数。那将是可以解决我的问题的UDAF;不幸的是我没有;没有。此外,结果预计将是熊猫的数据框。 StringType() 会起作用吗?
猜你喜欢
  • 1970-01-01
  • 2013-07-03
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2021-03-20
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多