【问题标题】:Flatmap on dataframe数据框上的平面图
【发布时间】:2017-10-15 03:47:09
【问题描述】:

在 spark 中对 DataFrame 执行 flatMap 的最佳方法是什么? 通过四处搜索并进行一些测试,我想出了两种不同的方法。这两个都有一些缺点,所以我认为应该有一些更好/更简单的方法来做到这一点。

我发现的第一种方法是先将DataFrame 转换为RDD,然后再返回:

val map = Map("a" -> List("c","d","e"), "b" -> List("f","g","h"))
val df = List(("a", 1.0), ("b", 2.0)).toDF("x", "y")

val rdd = df.rdd.flatMap{ row =>
  val x = row.getAs[String]("x")
  val x = row.getAs[Double]("y")
  for(v <- map(x)) yield Row(v,y)
}
val df2 = spark.createDataFrame(rdd, df.schema)

第二种方法是在使用flatMap之前创建一个DataSet(使用与上面相同的变量),然后再转换回来:

val ds = df.as[(String, Double)].flatMap{
  case (x, y) => for(v <- map(x)) yield (v,y)
}.toDF("x", "y")

当列数很少时,这两种方法都可以很好地工作,但是我的列数多于 2 列。有没有更好的方法来解决这个问题?最好采用不需要转换的方式。

【问题讨论】:

标签: scala apache-spark dataframe flatmap


【解决方案1】:

您可以从您的map RDD 创建第二个dataframe:

val mapDF = Map("a" -> List("c","d","e"), "b" -> List("f","g","h")).toList.toDF("key", "value")

然后执行join 并应用explode 函数:

val joinedDF = df.join(mapDF, df("x") === mapDF("key"), "inner")
  .select("value", "y")
  .withColumn("value", explode($"value"))

你得到了解决方案。

joinedDF.show()

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2019-09-08
    • 1970-01-01
    • 2019-11-27
    • 2020-11-17
    • 1970-01-01
    • 1970-01-01
    • 2021-05-16
    • 1970-01-01
    相关资源
    最近更新 更多