【问题标题】:How to split an RDD into different RDD's based on a value and give every part to a function如何根据一个值将 RDD 拆分为不同的 RDD,并将每个部分赋予一个函数
【发布时间】: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


【解决方案1】:

您真正需要做的就是通过聚合data 来计数,具体取决于布尔值可以采用的 2 个值。剩下的就是一个简单的计算,只依赖于这两个值。

val rdd = sc.parallelize(
  Seq(Element(true,"a"),Element(false,"a"),Element(true,"a"),
    Element(false,"b"),Element(false,"b"),Element(true,"b")))

val log2 = math.log(2)

// calculate an RDD[(String, (Int, Int))], first element of the tuple is the number of "true"s, and the second the number of "false"s
val entropy = rdd.map(e => (e.data, e.target)).aggregateByKey((0, 0))({
  case ((t, f), target) => if (target) (t + 1, f) else (t, f + 1)
}, {
  case ((t1, f1), (t2, f2)) => (t1 + t2, f1 + f2)
}).mapValues {
  case (t, f) =>
    val total = (t + f).toDouble
    val trueRatio = t.toDouble / total
    val falseRatio = f.toDouble / total
    -trueRatio * math.log(trueRatio) / log2 + falseRatio * math.log(falseRatio) / log2
}

// entropy is an RDD[(String, Double)]
entropy foreach println
// (a,-0.1383458330929479)
// (b,0.1383458330929479)

【讨论】:

    【解决方案2】:

    使用 DataFrame 的答案:

    import spark.implicits._
    
    val rdd = sc.parallelize(
          Seq(Element(true,"a"),Element(false,"a"),Element(true,"a"),
            Element(false,"b"),Element(false,"b"),Element(true,"b")))
    
    val log2 = math.log(2)
    
    var df = rdd.toDF()
    
    val groupedData = df.groupBy($"data")
      .agg(count(when($"target" === true, 1)).alias("true"),count(when($"target" === false, 1)).alias("false"))
      .withColumn("total", $"true" + $"false")
      .withColumn("true ratio", $"true" / $"total")
      .withColumn("false ratio", $"false" / $"total")
      .withColumn("entropy", -$"true ratio" * log($"true ratio") / log2 + $"false ratio" * log($"false ratio") / log2)
      .show()
    

    输出:

    +----+----+-----+-----+------------------+------------------+-------------------+
    |data|true|false|total|        true ratio|       false ratio|            entropy|
    +----+----+-----+-----+------------------+------------------+-------------------+
    |   b|   1|    2|    3|0.3333333333333333|0.6666666666666666| 0.1383458330929479|
    |   a|   2|    1|    3|0.6666666666666666|0.3333333333333333|-0.1383458330929479|
    +----+----+-----+-----+------------------+------------------+-------------------+
    

    【讨论】:

      猜你喜欢
      • 2016-01-03
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多