【发布时间】:2018-11-21 13:32:56
【问题描述】:
我对 spark 和大数据世界完全陌生。我有一个代码,它实际上创建了一个拆分 CSV 文件并返回两个字段的函数。
然后是 map 函数,我知道它是如何工作的,但我在代码的下一部分(操作发生在 totalsByAge 变量上)感到困惑,mapValues 和 reduceByKey 正在应用。请帮助我了解 reduceByKey 和 mapValues 在这里的工作原理?
def parseLine(line):
fields = line.split(',')
age = int(fields[2])
numFriends = int(fields[3])
return (age,numFriends)
line = sparkCont.textFile("D:\\ResearchInMotion\\ml-100k\\fakefriends.csv")
rdd = line.map(parseLine)
totalsByAge = rdd.mapValues(lambda x: (x, 1)).reduceByKey(lambda x, y: (x[0] + y[0], x[1] + y[1]))
averagesByAge = totalsByAge.mapValues(lambda x: x[0] / x[1])
results = averagesByAge.collect()
for result in results:
print(result)
我需要 totalsByAge 变量处理方面的帮助。如果您还可以详细说明对 averagesByAge 所做的操作,那将是很好的,如果缺少任何内容,请告诉我。 p>
【问题讨论】:
标签: python apache-spark pyspark rdd