Spark按段落计算单词平均长度问题求助
解决PySpark按段落计算单词平均长度的问题
嗨,我明白你现在的困扰啦!你的代码目前只是把每个段落的(单词长度, 1)元组列表拼在了一起,但没对总长度和单词数分别求和,所以没法直接算出平均值。咱们一步步修改代码,让它达到预期效果:
问题分析
原代码里的reduceByKey(lambda a,b: a+b)是把两个列表拼接起来,而不是对每个元组的两个元素分别累加。比如对于段落1,你得到的是一堆(长度,1)的元组,但我们需要的是(总长度, 总单词数),这样才能用总长度/总单词数算出平均。
修改后的完整代码
from pyspark import SparkContext, SparkConf # 初始化SparkContext conf = SparkConf().setAppName("paragraph_word_average") sc = SparkContext(conf=conf) # 读取文本文件 text = sc.textFile("walden.txt") # 1. 拆分段落ID和内容,过滤格式无效或内容为空的行 # 先split成[段落ID, 内容],确保格式正确且内容非空,避免后续报错 valid_lines = text.map(lambda line: line.split("|")).filter(lambda parts: len(parts) == 2 and parts[1].strip() != "") # 2. 把每个段落的内容拆分成单词,生成(段落ID, (单词长度, 1))的键值对 # flatMapValues会把每个单词单独拆出来,对应同一个段落ID word_metrics = valid_lines.flatMapValues(lambda content: content.split())\ .map(lambda kv: (kv[0], (len(kv[1]), 1))) # 3. 按段落ID聚合,累加总长度和总单词数 total_metrics = word_metrics.reduceByKey(lambda a, b: (a[0] + b[0], a[1] + b[1])) # 4. 计算每个段落的单词平均长度(处理除数为0的情况,避免报错) average_lengths = total_metrics.mapValues(lambda x: round(x[0]/x[1], 2) if x[1] != 0 else 0.0) # 保存结果 average_lengths.saveAsTextFile("results") # 停止SparkContext sc.stop()
关键步骤解释
- 拆分与过滤:用
split("|")直接拿到段落ID和内容,同时过滤掉格式错误(比如没有|)或内容为空的行,避免后续出现索引越界或无效计算。 - 生成单词指标:
flatMapValues把每个段落的单词拆成单独的条目,每个条目对应(段落ID, (单词长度, 1)),这样每个单词都能参与后续的累加计算。 - 聚合计算:
reduceByKey里的lambda函数把两个元组的第一个元素(长度)相加,第二个元素(计数)相加,最终得到每个段落的(总长度, 总单词数)。 - 计算平均值:用总长度除以总单词数,加上
round保留两位小数让结果更美观,同时处理总单词数为0的边界情况,避免出现除以0的报错。
这样修改后,你的输出就会变成类似('1', 4.2)、('2', 3.8)这样的格式,每个段落ID对应它的单词平均长度啦!
内容的提问来源于stack exchange,提问作者Neelhtak
相关产品推荐
相关产品推荐

