【发布时间】:2019-11-09 21:49:37
【问题描述】:
我正在尝试使用 PySpark 查找相邻元组列表之间的平均差异。
例如,如果我有这样的 RDD
vals = [(2,110),(2,130),(2,120),(3,200),(3,206),(3,206),(4,150),(4,160),(4,170)]
我想找出每个键的平均差异。
例如键值“2”
平均差异为 (abs(110-130) + abs(130-120))/2 = 15。
这是我目前的方法。我正在尝试更改平均计算代码以适应这种情况。但它似乎不起作用。
from pyspark import SparkContext
aTuple = (0,0)
interval = vals.aggregateByKey(aTuple, lambda a,b: (abs(a[0] - b),a[1] + 1),
lambda a,b: (a[0] + b[0], a[1] + b[1]))
finalResult = interval.mapValues(lambda v: (v[0]/v[1])).collect()
我想使用 RDD 函数来执行此操作,而不是使用 Spark SQL 或任何其他附加包。
最好的方法是什么?
如果您有任何问题,请告诉我。
感谢您的宝贵时间。
【问题讨论】:
标签: python apache-spark pyspark rdd moving-average