【问题标题】:How to count frequency of columns in row in a typedpipe in scalding?如何在烫伤中计算类型管道中行中列的频率?
【发布时间】:2016-05-22 22:19:34
【问题描述】:

我目前正在使用 scalding 进行 mapreduce 工作。我试图根据我在 typedpipe 中的行中看到特定值的次数来设定阈值。例如,如果我的 typedpipe 中有这些行:

第 1 列 |第 2 栏

'嗨' | '嘿'

'嗨' | '嗬'

'嗨' | '嗬'

'再见' | '再见'

我想将我在每行中看到的第 1 列和第 2 列中的值的频率附加到每一行。意思是输出看起来像:

第 1 列 |第 2 栏 |第 1 列频率 |第 2 列频率

'嗨' | '嘿'| 3 | 1

'嗨' | '嗬' | 3 | 2

'嗨' | '嗬' | 3 | 2

'再见' | '再见' | 1 | 1

目前,我通过按每列对类型化的管道进行分组来做到这一点,如下所示:

  val key2Freqs = input.groupBy('key2) {
    _.size('key2Freq)
  }.rename('key2 -> 'key2Right).project('key2Right, 'key2Freq);

然后像这样使用 key2Freqs 加入原始输入:

  .joinWithSmaller('key2 -> 'key2Right, key2Freqs, joiner = new LeftJoin)

但是,这真的很慢,而且在我看来,对于本质上非常简单的任务来说效率很低。它变得特别长 b/c 我有 6 个不同的键,我想获得这些值,我目前正在映射和加入 6 个不同的工作时间。一定有更好的方法来做到这一点,对吧?

【问题讨论】:

    标签: scala hadoop mapreduce scalding


    【解决方案1】:

    如果每列中不同值的数量足够小以将它们全部放入内存中,您可以将您的列 .map 放入 Map[String,Int],然后 .groupAll.sum 一次性将它们全部计数(我是使用“typed api”表示法,不太记得在字段 api 中这是如何完成的,但你明白了)。您需要使用来自algebirdMapMonoid,或者如果您不想为此添加依赖项,则只需编写自己的,这并不难。 然后,您将得到一个管道,其中包含生成的Map 的单个条目。现在,您可以获得原始管道,然后执行 .crossWithTiny 将带有计数的地图带入其中,然后执行 .map 以提取单个计数。

    否则,如果您无法将所有这些都保存在内存中,那么您现在正在做什么似乎是唯一的方法......除非您实际上是在寻找“顶级击球手”的近似值,而不是确切的计数整个宇宙的...在这种情况下,请查看 algebird 的SketchMap

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2014-06-19
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2020-09-06
      • 2017-07-13
      相关资源
      最近更新 更多