【发布时间】:2020-04-24 01:16:28
【问题描述】:
我有一个 RDD,其中每个元素都是一个案例类,如下所示:
case class Element(target: Boolean, data: String)
现在我需要根据字符串数据(它是一个离散变量)来拆分 RDD。
然后在每次拆分时执行一个函数def f(elements: RDD[Element]): Double。
我试图像这样创建一个 pairRDD:val test = elementsRDD.map(E => (E.data, E)) 所以我有 (key, value) 对,但我不知道在此之后该怎么做(如何拆分它们,因为 groupBy 会返回 Iteravle(V) 和不是所有值的 RDD)。
我还可以过滤data: String 的每个可能值并对结果执行函数 f。但我不知道 ´´´data: String´´´ 可以提前采用的所有可能值。而且,首先检查所有数据以检查不同的可能性,然后再对其进行多次过滤似乎效率不高。
那么有没有什么方法可以高效完成呢?
【问题讨论】:
-
我并没有真正看到与首先获取所有不同的
data值,然后过滤原始RDD 以创建N不同的RDD 不同的方法。问题是f将 RDD 作为参数。我们能否更深入地了解f的作用?比如f(elements) = elements.count(),通过简单的聚合就可以轻松解决问题。 -
f 计算熵,因此它计算具有 target = true 和 target = false 的元素的数量和总数。然后它返回:-((AmountTrue/total)*log2(AmountTrue/total) + ((AmountFalse/total)*log2(AmountFalse/total)
标签: scala apache-spark rdd