在PySpark UDF中使用外部库引发Pickle错误的解决方法
PySpark UDF 结合 pymorphy2 报错的最简解决方法
问题代码
import pandas as pd from pymorphy2 import MorphAnalyzer from pyspark.sql import SparkSession from pyspark.sql import types as T from pyspark.sql import functions as F spark = SparkSession.builder.appName("udf").getOrCreate() def gender(s): m = MorphAnalyzer() return m.parse(s)[0].tag.gender gen = F.udf(gender, T.StringType()) df = spark.createDataFrame(pd.DataFrame({"name": ["кирилл", "вавила"]})) df.select(gen("name").alias("gender")).show()
报错信息
ERROR Executor: Exception in task 2.0 in stage 29.0 (TID 151) net.razorvine.pickle.PickleException: expected zero arguments for construction of ClassDict (for pyspark.cloudpickle.cloudpickle._make_skeleton_class). This happens when an unsupported/unregistered class is being unpickled that requires construction arguments. Fix it by registering a custom IObjectConstructor for this class. at net.razorvine.pickle.objects.ClassDictConstructor.construct(ClassDictConstructor.java:23) at net.razorvine.pickle.Unpickler.load_reduce(Unpickler.java:759) at net.razorvine.pickle.Unpickler.dispatch(Unpickler.java:199)
最简解决方法
将MorphAnalyzer的实例化逻辑移到UDF函数外部,让每个Executor进程仅初始化一次该实例,即可规避错误:
import pandas as pd from pymorphy2 import MorphAnalyzer from pyspark.sql import SparkSession from pyspark.sql import types as T from pyspark.sql import functions as F spark = SparkSession.builder.appName("udf").getOrCreate() # 移至UDF外部,每个Executor进程仅初始化一次 m = MorphAnalyzer() def gender(s): return m.parse(s)[0].tag.gender gen = F.udf(gender, T.StringType()) df = spark.createDataFrame(pd.DataFrame({"name": ["кирилл", "вавила"]})) df.select(gen("name").alias("gender")).show()
原因说明
原代码在每次调用UDF时都会创建新的MorphAnalyzer实例,该实例在序列化发送到Executor节点时,会因为pymorphy2内部对象无法被Spark序列化器正确处理而触发pickle错误。将实例化逻辑移到UDF外部后,每个Executor进程只会初始化一次实例,既解决了序列化问题,也提升了性能(无需重复创建实例)。
内容的提问来源于stack exchange,提问作者Sergey Bushmanov
相关产品推荐
相关产品推荐

