PySpark:如何用Lambda操作生成相邻词对计数元组?
问题描述
已通过PySpark处理得到单个单词组成的RDD,输出示例:
['vita', 'oscura', 'smarrita', 'dura', 'forte', 'paura', 'morte', 'trovai', 'scorte', 'v’intrai']
需要将其转换为相邻词对+计数1的元组列表,格式示例:
[('vita','oscura',1),('oscura','smarrita',1),('smarrita','dura',1), ('dura','forte',1) ...]
尝试用Lambda函数实现时出现索引越界错误,现有代码如下:
def lower_clean_str(x): punc='!\"#$%&\'()*+,-./:;<=>?@[\\]^_`{|}~' lowercased_str = x.lower() for ch in punc: lowercased_str = lowercased_str.replace(ch, '') return lowercased_str clean_dcr=dcr.map(lower_clean_str) print(clean_dcr.take(10)) # 按空格拆分后取最后一个单词 clean_dcr=clean_dcr.map(lambda line: line.split()[-1]) print(clean_dcr.take(10)) # 以下代码报错 # clean_dcr=clean_dcr.map((lambda line:line[0][0],line[0][1])),1) # print(clean_dcr.take(3))
解决方法
错误原因
当前clean_dcr的每个元素是单个字符串单词(比如"vita"),不是数组/列表,你尝试用line[0][0]去索引单个字符串的子元素,必然会出现索引越界。
正确实现方式
要生成相邻词对,需要将RDD中的单词按顺序两两配对,这里提供两种可行方案:
方案1:使用Zip操作
将原RDD与去掉第一个元素的RDD进行Zip,得到相邻词对后再添加计数1:
# 生成去掉第一个元素的RDD rdd_tail = clean_dcr.zipWithIndex().filter(lambda x: x[1] > 0).map(lambda x: x[0]) # 配对并添加计数 pair_rdd = clean_dcr.zip(rdd_tail).map(lambda x: (x[0], x[1], 1)) print(pair_rdd.take(5))
方案2:使用Sliding窗口(Spark 2.4+支持)
利用RDD的sliding方法直接生成长度为2的滑动窗口,再转换为目标格式:
# 生成滑动窗口,每个窗口包含2个相邻单词 pair_rdd = clean_dcr.sliding(2).map(lambda window: (window[0], window[1], 1)) print(pair_rdd.take(5))
代码替换说明
直接替换你代码中报错的部分即可,两种方案都能输出符合要求的(word1, word2, 1)格式元组。
内容的提问来源于stack exchange,提问作者jeannetton
相关产品推荐
相关产品推荐

