【问题标题】:PicklingError: Could not serialize object: TypeError: cannot serialize '_io.BufferedReader' objectPicklingError:无法序列化对象:TypeError:无法序列化“_io.BufferedReader”对象
【发布时间】:2020-08-04 01:20:19
【问题描述】:

我正在基于对文本进行词形还原来实现情感分析(在调用 udf 之前,我正在根据词形还原函数的要求将格式转换为字符串),但我遇到了一个我无法理解的酸洗错误。

PicklingError:无法序列化对象:TypeError:无法序列化“_io.BufferedReader”对象

下面是我的代码

    def get_part_of_speech(word):
    probable_part_of_speech = wordnet.synsets(word)
  
    pos_counts = Counter()

    pos_counts["n"] = len(  [ item for item in probable_part_of_speech if item.pos()=="n"]  )
    pos_counts["v"] = len(  [ item for item in probable_part_of_speech if item.pos()=="v"]  )
    pos_counts["a"] = len(  [ item for item in probable_part_of_speech if item.pos()=="a"]  )
    pos_counts["r"] = len(  [ item for item in probable_part_of_speech if item.pos()=="r"]  )
      
    most_likely_part_of_speech = pos_counts.most_common(1)[0][0]
    return most_likely_part_of_speech


def Lemmatizing_Words(Words):
    Lm = WordNetLemmatizer()
    Lemmatized_Words = []
    for word in Words:
        Lemmatized_Words.append(Lm.lemmatize(word,get_part_of_speech(word)))
    return Lemmatized_Words

    sparkLemmer1 = udf(lambda x:Lemmatizing_Words(x), StringType())
    df_new_lemma_1 = dfStopwordRemoved.select('overall',sparkLemmer1('filteredreviewText'))

我面临的错误:

  PicklingError       Traceback (most recent call last)
<ipython-input-20-1a663a7b63cb> in <module>()
---->1 df_new_lemma_1 = dfStopwordRemoved.select('overall',sparkLemmer1('filteredreviewText'))


~/.conda/envs/myenv/lib/python3.5/site-packages/pyspark/sql/udf.py in wrapper(*args)
    187         @functools.wraps(self.func, assigned=assignments)
    188         def wrapper(*args):
--> 189             return self(*args)
    190 
    191         wrapper.__name__ = self._name

~/.conda/envs/myenv/lib/python3.5/site-packages/pyspark/sql/udf.py in __call__(self, *cols)
    165 
    166     def __call__(self, *cols):
--> 167         judf = self._judf
    168         sc = SparkContext._active_spark_context
    169         return Column(judf.apply(_to_seq(sc, cols, _to_java_column)))

~/.conda/envs/myenv/lib/python3.5/site-packages/pyspark/sql/udf.py in _judf(self)
    149         # and should have a minimal performance impact.
    150         if self._judf_placeholder is None:
--> 151             self._judf_placeholder = self._create_judf()
    152         return self._judf_placeholder
    153 

~/.conda/envs/myenv/lib/python3.5/site-packages/pyspark/sql/udf.py in _create_judf(self)
    158         sc = spark.sparkContext
    159 
--> 160         wrapped_func = _wrap_function(sc, self.func, self.returnType)
    161         jdt = spark._jsparkSession.parseDataType(self.returnType.json())
    162         judf = sc._jvm.org.apache.spark.sql.execution.python.UserDefinedPythonFunction(

~/.conda/envs/myenv/lib/python3.5/site-packages/pyspark/sql/udf.py in _wrap_function(sc, func, returnType)
     33 def _wrap_function(sc, func, returnType):
     34     command = (func, returnType)
---> 35     pickled_command, broadcast_vars, env, includes = _prepare_for_python_RDD(sc, command)
     36     return sc._jvm.PythonFunction(bytearray(pickled_command), env, includes, sc.pythonExec,
     37                                   sc.pythonVer, broadcast_vars, sc._javaAccumulator)

~/.conda/envs/myenv/lib/python3.5/site-packages/pyspark/rdd.py in _prepare_for_python_RDD(sc, command)
   2418     # the serialized command will be compressed by broadcast
   2419     ser = CloudPickleSerializer()
-> 2420     pickled_command = ser.dumps(command)
   2421     if len(pickled_command) > (1 << 20):  # 1M
   2422         # The broadcast will have same life cycle as created PythonRDD

~/.conda/envs/myenv/lib/python3.5/site-packages/pyspark/serializers.py in dumps(self, obj)
    605                 msg = "Could not serialize object: %s: %s" % (e.__class__.__name__, emsg)
    606             cloudpickle.print_exec(sys.stderr)
--> 607             raise pickle.PicklingError(msg)
    608 
    609 


PicklingError: Could not serialize object: TypeError: cannot serialize '_io.BufferedReader' object

【问题讨论】:

    标签: python pyspark nltk pickle wordnet


    【解决方案1】:

    这是pickle 模块的一个常见问题和缺陷,它无法持久化某些类型的对象,并且 PySpark 正在为进程间通信腌制对象以分配计算。

    此外,腌制已打开文件或流的句柄是没有意义的(NLTK 中的WordNetLematizer 使用秘密的_io.BufferedReader 句柄指向 WordNet 文件。

    要解决此问题,您需要复制 NLTK 的 WordNetLematizer 的实现并预先加载数据以在 udf PySpark 的函数中使用 WordNetLematizer。

    【讨论】:

      猜你喜欢
      • 2020-11-25
      • 2021-03-20
      • 2018-09-18
      • 1970-01-01
      • 2018-07-23
      • 1970-01-01
      • 1970-01-01
      • 2018-04-25
      • 1970-01-01
      相关资源
      最近更新 更多