【发布时间】: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