【问题标题】:Functional Programming way to calculate something like a rolling sum计算滚动和之类的函数式编程方法
【发布时间】:2017-06-16 17:35:41
【问题描述】:

假设我有一个数字列表:

val list = List(4,12,3,6,9)

对于列表中的每个元素,我需要找到滚动总和,即最终输出应该是:

List(4, 16, 19, 25, 34)

是否有任何转换允许我们将列表的两个元素(当前和前一个)作为输入并基于两者进行计算? 类似map(initial)((curr,prev) => curr+prev)

我想在不维护任何共享全局状态的情况下实现这一目标。

编辑:我希望能够在 RDD 上进行相同类型的计算。

【问题讨论】:

    标签: scala functional-programming rolling-sum


    【解决方案1】:

    您可以使用scanLeft

    list.scanLeft(0)(_ + _).tail
    

    【讨论】:

    • 谢谢!这个有效。但是,我想在 spark RDD 上实现相同的目标,而 scanLeft 似乎没有在 RDD 上实现。是否也可以对 RDD 进行类似操作?
    • 这里是scanLeft对于RDD的描述和实现:erikerlandson.github.io/blog/2014/08/09/…
    【解决方案2】:

    下面的cumSum 方法应该适用于任何RDD[N],其中N 有一个隐式Numeric[N] 可用,例如IntLongBigIntDouble

    import scala.reflect.ClassTag
    import org.apache.spark.rdd.RDD
    
    def cumSum[N : Numeric : ClassTag](rdd: RDD[N]): RDD[N] = {
      val num = implicitly[Numeric[N]]
      val nPartitions = rdd.partitions.length
    
      val partitionCumSums = rdd.mapPartitionsWithIndex((index, iter) => 
        if (index == nPartitions - 1) Iterator.empty
        else Iterator.single(iter.foldLeft(num.zero)(num.plus))
      ).collect
       .scanLeft(num.zero)(num.plus)
    
      rdd.mapPartitionsWithIndex((index, iter) => 
        if (iter.isEmpty) iter
        else {
          val start = num.plus(partitionCumSums(index), iter.next)
          iter.scanLeft(start)(num.plus)
        }
      )
    }
    

    将这种方法推广到任何具有“零”(即任何幺半群)的关联二元运算符应该是相当简单的。关联性是并行化的关键。如果没有这种关联性,您通常会被困在以串行方式遍历RDD 的条目。

    【讨论】:

    • 好吧,我试过这个:cumSum(sc.parallelize(List(1,4,2,5,3))) 但得到了一个例外java.lang.UnsupportedOperationException: empty.reduceLeft。看起来它试图减少一个空集合
    • 这可能是因为您的RDD 的某些分区是空的。我稍微修改了代码。希望它现在可以工作。
    • 这是一个非常有用的算法!谢谢@Jason
    【解决方案3】:

    我不知道 spark RDD 支持哪些功能,所以我不确定这是否满足您的条件,因为我不知道是否支持 zipWithIndex(如果答案没有帮助,请通过评论,我将删除我的答案):

    list.zipWithIndex.map{x => list.take(x._2+1).sum}
    

    这段代码对我有用,它总结了元素。它获取列表元素的索引,然后将相应的 n 个元素添加到列表中(注意 +1,因为 zipWithIndex 以 0 开头)。

    打印时,我得到以下信息:

    List(4, 16, 19, 25, 34)
    

    【讨论】:

    • 您确定这适用于并行集合吗? parList.take(n) 总是返回前 n 个元素吗?还是会随机返回 n 个元素?此外,这在计算上会很昂贵,因为您基本上是通过蛮力为循环中的每次迭代添加前 n 个元素。在我看来,函数式编程的误用。
    • 这是一个 O(n) 问题,但您使用的是 O(n^2) 算法。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2015-01-15
    • 1970-01-01
    • 2017-05-13
    • 2010-12-04
    • 1970-01-01
    • 2011-01-16
    相关资源
    最近更新 更多