【发布时间】:2016-08-02 16:39:15
【问题描述】:
所以这个标题应该足够令人困惑,所以我会尽力解释。我正在尝试将此函数分解为已定义的函数,以便更好地了解 aggregateByKey 如何为将要写入我的代码的其他团队工作。我有以下聚合:
val firstLetter = stringRDD.aggregateByKey(Map[Char, Int]())(
(accumCount, value) => accumCount.get(value.head) match {
case None => accumCount + (value.head -> 1)
case Some(count) => accumCount + (value.head -> (count + 1))
},
(accum1, accum2) => accum1 ++ accum2.map{case(k,v) => k -> (v + accum1.getOrElse(k, 0))}
).collect()
我一直想把它分解如下:
val firstLet = Map[Char, Int]
def fSeq(accumCount:?, value:?) = {
accumCount.get(value.head) match {
case None => accumCount + (value.head -> 1)
case Some(count) => accumCount + (value.head -> (count + 1))
}
}
def fComb(accum1:?, accum2:?) = {
accum1 ++ accum2.map{case(k,v) => k -> (v + accum1.getOrElse(k, 0))
}
由于初始值是 Map[Char, Int] 我不确定要定义什么 accumCount, Value 数据类型。我尝试了不同的东西,但似乎没有任何效果。有人可以帮我定义数据类型并解释你是如何确定的吗?
【问题讨论】:
-
这里的输入是什么?
RDD[(T, String)]?
标签: scala apache-spark