【问题标题】:How to Reduce by key in "Scala" [Not In Spark]如何在“Scala”中按键减少 [不在 Spark 中]
【发布时间】:2019-02-09 08:10:32
【问题描述】:

我正在尝试减少 Scala 中的键值,是否有任何方法可以根据 Scala 中的键来减少值。 [我知道我们可以通过 spark 中的 reduceByKey 方法来做,但是我们如何在 Scala 中做同样的事情? ]

输入数据是:

val File = Source.fromFile("C:/Users/svk12/git/data/retail_db/order_items/part-00000")
                 .getLines()
                 .toList

 val map = File.map(x => x.split(","))
               .map(x => (x(1),x(4)))

  map.take(10).foreach(println)

在上述步骤之后,我得到的结果是:

(2,250.0)
(2,129.99)
(4,49.98)
(4,299.95)
(4,150.0)
(4,199.92)
(5,299.98)
(5,299.95)

预期结果:

(2,379.99)
(5,499.93)
.......

【问题讨论】:

  • 我认为这缺少一个分组步骤。

标签: scala higher-order-functions


【解决方案1】:

Scala 2.13 开始,您可以使用groupMapReduce 方法(顾名思义)等效于groupBy,后跟mapValuesreduce 步骤:

io.Source.fromFile("file.txt")
  .getLines.to(LazyList)
  .map(_.split(','))
  .groupMapReduce(_(1))(_(4).toDouble)(_ + _)

groupMapReduce 舞台:

  • groups 按第二个元素 (_(1)) 拆分数组(groupMapReduce 的组部分)

  • maps 每个组中的每个数组出现到其第 4 个元素,并将其转换为 Double (_(4).toDouble)(映射组的一部分MapReduce)

  • 每个组 (_ + _) 中的reduces 值通过求和(减少 groupMap 的一部分Reduce)。

这是one-pass version 可以翻译的内容:

seq.groupBy(_(1)).mapValues(_.map(_(4).toDouble).reduce(_ + _))

还要注意从IteratorLazyList 的转换,以便使用提供groupMapReduce 的集合(我们不使用Stream,因为从Scala 2.13 开始,建议替换LazyList Streams)。

【讨论】:

    【解决方案2】:

    您似乎想要文件中某些值的总和。一个问题是文件是字符串,因此您必须将String 转换为数字格式才能对其求和。

    这些是您可能会使用的步骤。

    io.Source.fromFile("so.txt") //open file
      .getLines()                //read line-by-line
      .map(_.split(","))         //each line is Array[String]
      .toSeq                     //to something that can groupBy()
      .groupBy(_(1))             //now is Map[String,Array[String]]
      .mapValues(_.map(_(4).toInt).sum) //now is Map[String,Int]
      .toSeq                     //un-Map it to (String,Int) tuples
      .sorted                    //presentation order
      .take(10)                  //sample
      .foreach(println)          //report
    

    如果任何文件数据不符合要求的格式,这当然会抛出。

    【讨论】:

      【解决方案3】:

      没有内置任何东西,但你可以这样写:

      def reduceByKey[A, B](items: Traversable[(A, B)])(f: (B, B) => B): Map[A, B] = {
        var result = Map.empty[A, B]
        items.foreach {
          case (a, b) =>
            result += (a -> result.get(a).map(b1 => f(b1, b)).getOrElse(b))
        }
        result
      }
      

      有一些空间可以优化这一点(例如,使用可变映射),但总体思路保持不变。

      另一种方法,更具声明性但效率较低(创建多个中间集合;可以重写但不清晰:

      def reduceByKey[A, B](items: Traversable[(A, B)])(f: (B, B) => B): Map[A, B] = {
        items
          .groupBy { case (a, _) => a }
          .mapValues(_.map { case (_, b) => b }.reduce(f))
          // mapValues returns a view, view.force changes it back to a realized map
          .view.force
      }
      

      【讨论】:

      • 最好对这个函数进行柯里化,以便编译器知道B 的类型对于f 参数。此外,fold 优于 map/getOrElse
      • 是的,currying 确实是有道理的,但 fold 是一个品味问题 :) 我个人不喜欢它看起来像 fold(a)(_ + b) 而不是 fold(none = a, some = _ + b) (例如 scalaz 中的 cata) .
      【解决方案4】:

      首先使用key对元组进行分组,这里是第一个元素,然后是reduce。 以下代码将起作用 -

      val reducedList = map.groupBy(_._1).map(l => (l._1, l._2.map(_._2).reduce(_+_)))
      print(reducedList)
      

      【讨论】:

        【解决方案5】:

        这里使用 foldLeft 的另一种解决方案:

        val File : List[String] = ???
        
        File.map(x => x.split(","))
          .map(x => (x(1),x(4).toInt))
          .foldLeft(Map.empty[String,Int]){case (state, (key,value)) => state.updated(key,state.get(key).getOrElse(0)+value)}
          .toSeq
          .sortBy(_._1)
          .take(10)
          .foreach(println)
        

        【讨论】:

          猜你喜欢
          • 1970-01-01
          • 1970-01-01
          • 1970-01-01
          • 1970-01-01
          • 2016-12-05
          • 1970-01-01
          • 2018-05-08
          • 1970-01-01
          • 2015-07-08
          相关资源
          最近更新 更多