【发布时间】:2018-09-28 05:20:06
【问题描述】:
我正在使用 pyspark 流从 tweepy 收集数据。完成所有设置后,我通过 elasticsearch.index() 将 dict(json) 发送到 elasticsearch。但我得到“can't pickle_thread.lock objects”错误和其他 63 个错误。回溯日志太长,无法在我的控制台中显示!
设计是我得到一个 json/dict 类型的文件,将其转换为 DStream,通过在 map() 函数中调用 TextBlob 为其添加另一个特征名称“sentiment”。一切正常,但是当我添加另一个地图函数来调用 elasticsearch.index() 时,我得到了错误。
下面是我控制台中超长错误日志的一部分。
块引用 在处理上述异常的过程中,又出现了一个异常: 回溯(最近一次通话最后): 转储中的文件“/Users/ayane/anaconda/lib/python3.6/site-packages/pyspark/streaming/util.py”,第 105 行 func.func, func.rdd_wrap_func, func.deserializers))) 转储中的文件“/Users/ayane/anaconda/lib/python3.6/site-packages/pyspark/serializers.py”,第 460 行 返回 cloudpickle.dumps(obj, 2) 转储中的文件“/Users/ayane/anaconda/lib/python3.6/site-packages/pyspark/cloudpickle.py”,第 704 行 cp.dump(obj) 转储中的文件“/Users/ayane/anaconda/lib/python3.6/site-packages/pyspark/cloudpickle.py”,第 162 行 提出 pickle.PicklingError(msg) _pickle.PicklingError:无法序列化对象:TypeError:无法腌制 _thread.lock 对象 在 org.apache.spark.streaming.api.python.PythonTransformFunctionSerializer$.serialize(PythonDStream.scala:144) 在 org.apache.spark.streaming.api.python.TransformFunction$$anonfun$writeObject$1.apply$mcV$sp(PythonDStream.scala:101) 在 org.apache.spark.streaming.api.python.TransformFunction$$anonfun$writeObject$1.apply(PythonDStream.scala:100) 在 org.apache.spark.streaming.api.python.TransformFunction$$anonfun$writeObject$1.apply(PythonDStream.scala:100) 在 org.apache.spark.util.Utils$.tryOrIOException(Utils.scala:1303) ... 63 更多
我的部分代码如下所示:
def sendPut(doc):
res = es.index(index = "tweetrepository", doc_type= 'tweet', body = doc)
return doc
myJson = dataStream.map(decodeJson).map(addSentiment).map(sendPut)
myJson.pprint()
这里是decodeJson函数:
def decodeJson(str):
return json.loads(str)
这里是 addSentiment 函数:
def addSentiment(dic):
dic['Sentiment'] = get_tweet_sentiment(dic['Text'])
return dic
这里是 get_tweet_sentiment 函数:
def get_tweet_sentiment(tweet):
analysis = TextBlob(tweet)
if analysis.sentiment.polarity > 0:
return 'positive'
elif analysis.sentiment.polarity == 0:
return 'neutral'
else:
return 'negative'
【问题讨论】:
-
pprint() 的输出在哪里?
-
如果我删除最后一个 map(sendPut) 函数,输出将是: {"Text": "Hello from the XXXXX", "Loc": "123.00, 25.36", "Timestamp" : "20180418125154", "情绪": "积极"}
-
需要看类型,因为有些数据类型是不可pickleable的对象。向我们提供更多信息,并分离出您的映射函数。
-
TextBlob 是一个用于 NLP 的 API。
-
我知道 TextBlob 是什么,它是一个不可腌制的对象。分离出你的映射并告诉我你是否在 SendPut 中遇到错误?
标签: python apache-spark elasticsearch pyspark