Pyspark 如何从RDD键值对的value词列表中移除指定停用词
Spark RDD 停用词过滤正确实现
错误写法的核心问题
- 你代码中
stopwords是分布式RDD对象,不能直接在Executor端运行的lambda函数中,用in关键词判断元素是否属于RDD,这类操作只能针对本地集合执行 - 你需要操作的是pair RDD中每个键对应的词列表内部的元素,不是过滤整个键值对,也不是返回布尔类型的判断结果
正确实现步骤
首先需要把小体量的停用词收集为Driver端的本地集合,再通过广播变量分发到所有Executor节点,避免重复传输占用资源,再对词列表做过滤即可,代码如下:
# 将停用词RDD转为本地集合后广播 stopword_set = set(stopwords.collect()) stopword_broadcast = sc.broadcast(stopword_set) # 仅处理pair RDD的值部分,过滤掉停用词 pair1 = pair.mapValues(lambda words: [w for w in words if w not in stopword_broadcast.value])
结果验证
执行pair1.collect()后得到的结果示例如下:
[ ['2003', ['old', 'men']], ['2004', ['I', 'like', 'way', 'and', 'teachers', 'work']], ['2005', ['another', 'title', 'goes', 'here', 'for', 'any', 'reason']], ['2006', ['text', 'and', 'strings', 'are', 'similar']], ['2007', ['cowbows', 'love', 'to', 'ride']] ]
内容的提问来源于stack exchange,提问作者JohnDoe34
相关产品推荐
相关产品推荐

