PySpark RDD如何对字典列表执行reduceByKey完成词对计数合并
PySpark合并词对计数字典解决方案
核心实现两步走:
- 用
flatMap将每个字典的所有(词对, 计数)键值对展开为独立RDD元素 - 用
reduceByKey对相同词对的计数做累加
完整可运行示例代码:
from pyspark import SparkContext sc = SparkContext("local", "word_pair_count") # 模拟给出的输入文本RDD input_rdd = sc.parallelize([ ' the adventure of the blue carbuncle the adventure of the blue carbuncle the adventure of the blue carbuncle ', ' the adventure of the blue carbuncle' ]) # 词对字典生成逻辑,可替换为你已有的映射函数 def generate_pair_count(text): words = text.strip().split() pair_count = {} for i in range(len(words)-1): pair = (words[i], words[i+1]) pair_count[pair] = pair_count.get(pair, 0) + 1 return pair_count # 生成每个元素为单文本词对计数字典的中间RDD pair_dict_rdd = input_rdd.map(generate_pair_count) # 此时pair_dict_rdd第一个元素含('of','blue'):3,第二个元素含('of','blue'):1 # 核心合并逻辑 result_rdd = pair_dict_rdd \ .flatMap(lambda dict_item: dict_item.items()) \ .reduceByKey(lambda a, b: a + b) # 输出汇总结果 print(result_rdd.collectAsMap())
运行后可直接得到('of', 'blue'): 4的预期结果,所有词对的计数都会同步完成汇总。
如果之前的写法没有生效,大概率是flatMap阶段没有调用字典的items()方法,直接返回字典本身只会展开字典的键,丢失计数值,无法完成后续聚合。
内容的提问来源于stack exchange,提问作者Teddy
相关产品推荐
相关产品推荐

