【问题标题】:Homemade DataFrame aggregation/dropDuplicates Spark自制 DataFrame 聚合/dropDuplicates Spark
【发布时间】:2018-03-07 16:41:16
【问题描述】:

我想对我的 DataFrame df 执行转换,这样我每个键在最终 DataFrame 中只有一次且只有一次。

出于机器学习的目的,我不想在我的数据集中存在偏见。这绝不应该发生,但我从数据源获得的数据包含这种“怪异”。因此,如果我有具有相同键的行,我希望能够选择两者的组合(如平均值)或字符串连接(例如标签)或随机值集。

假设我的 DataFrame df 看起来像这样:

+---+----+-----------+---------+
|ID1| ID2|       VAL1|     VAL2|
+---+----+-----------+---------+
|  A|   U|     PIERRE|        1|
|  A|   U|     THOMAS|        2|
|  A|   U|    MICHAEL|        3|
|  A|   V|        TOM|        2|
|  A|   V|       JACK|        3|
|  A|   W|     MICHEL|        2|
|  A|   W|     JULIEN|        3|
+---+----+-----------+---------+

我希望我的最终 DataFrame out 仅随机保留每个键的一组值。它可能是另一种类型的聚合(例如将所有值串联为字符串),但我只是不想从中构建 Integer 值,而是构建新条目。

例如。最终输出可能是(仅保留每个键的第一行):

+---+----+-----------+---------+
|ID1| ID2|       VAL1|     VAL2|
+---+----+-----------+---------+
|  A|   U|     PIERRE|        1|
|  A|   V|        TOM|        2|
|  A|   W|     MICHEL|        2|
+---+----+-----------+---------+

另一个最终输出可能是(每个键保持随机行):

+---+----+-----------+---------+
|ID1| ID2|       VAL1|     VAL2|
+---+----+-----------+---------+
|  A|   U|    MICHAEL|        3|
|  A|   V|       JACK|        3|
|  A|   W|     MICHEL|        2|
+---+----+-----------+---------+

或者,建立一组新的价值观:

+---+----+--------------------------+----------+
|ID1| ID2|                      VAL1|      VAL2|
+---+----+--------------------------+----------+
|  A|   U| (PIERRE, THOMAS, MICHAEL)| (1, 2, 3)|
|  A|   V|               (TOM, JACK)|    (2, 3)|
|  A|   W|          (MICHEL, JULIEN)|    (2, 3)|
+---+----+--------------------------+----------+

答案应该是使用 Spark 和 Scala。我还想强调,实际的架构比这要复杂得多,我想找到一个通用的解决方案。此外,我不想仅从一列中获取唯一值,但过滤掉具有相同键的行。谢谢!

编辑这是我尝试做的(但Row.get(colname) 抛出NoSuchElementException: key not found...):

  def myDropDuplicatesRandom(df: DataFrame, colnames: Seq[String]): DataFrame = {
    val fields_map: Map[String, (Int, DataType)] =
      df.schema.fieldNames.map(fname => {
        val findex = df.schema.fieldIndex(fname)
        val ftype = df.schema.fields(findex).dataType
        (fname, (findex, ftype))
      }).toMap[String, (Int, DataType)]

    df.sparkSession.createDataFrame(
      df.rdd
        .map[(String, Row)](r => (colnames.map(colname => r.get(fields_map(colname)._1).toString.replace("`", "")).reduceLeft((x, y) => "" + x + y), r))
        .groupByKey()
        .map{case (x: String, y: Iterable[Row]) => Utils.randomElement(y)}
    , df.schema)
  }

【问题讨论】:

  • @David 我不是在一列上寻找不同的值,而是在寻找一种方法来过滤掉具有相同键的值。
  • 您似乎想根据 ID2(可能还有 ID1)删除重复项。不确定我是否理解您的问题与此有何不同。 “钥匙”是什么意思?
  • 因此,您可以根据ID1ID2 使用dropDuplicates
  • 啊,我的错。很难找到一个好的 scala spark 答案,我是 pyspark 用户。你会想要使用dropDuplicates。表单将类似于val out = df.dropDuplicates(Seq("ID1", "ID2"))

标签: scala apache-spark spark-dataframe rdd


【解决方案1】:

这是一种方法:

val df = Seq(
  ("A", "U", "PIERRE", 1),
  ("A", "U", "THOMAS", 2),
  ("A", "U", "MICHAEL", 3),
  ("A", "V", "TOM", 2),
  ("A", "V", "JACK", 3),
  ("A", "W", "MICHEL", 2),
  ("A", "W", "JULIEN", 3)
).toDF("ID1", "ID2", "VAL1", "VAL2")

import org.apache.spark.sql.functions._

// Gather key/value column lists based on specific filtering criteria
val keyCols = df.columns.filter(_.startsWith("ID"))
val valCols = df.columns diff keyCols

// Group by keys to aggregate combined value-columns then re-expand
df.groupBy(keyCols.map(col): _*).
  agg(first(struct(valCols.map(col): _*)).as("VALS")).
  select($"ID1", $"ID2", $"VALS.*")

// +---+---+------+----+
// |ID1|ID2|  VAL1|VAL2|
// +---+---+------+----+
// |  A|  W|MICHEL|   2|
// |  A|  V|   TOM|   2|
// |  A|  U|PIERRE|   1|
// +---+---+------+----+

[更新]

如果我正确理解您的扩展要求,您正在寻找一种通用方法来通过任意agg 函数的键转换数据帧,例如:

import org.apache.spark.sql.Column

def customAgg(keyCols: Seq[String], valCols: Seq[String], aggFcn: Column => Column) = {
  df.groupBy(keyCols.map(col): _*).
    agg(aggFcn(struct(valCols.map(col): _*)).as("VALS")).
    select($"ID1", $"ID2", $"VALS.*")
}

customAgg(keyCols, valCols, first)

我想说,沿着这条路走下去会导致非常有限的适用agg 函数。虽然以上适用于first,但您必须针对collect_list/collect_set 等进行不同的实现。当然可以手动滚动所有各种类型的agg 函数,但这可能会导致不必要的代码维护麻烦。

【讨论】:

    【解决方案2】:

    您可以将groupByfirststruct 一起使用,如下所示

      import org.apache.spark.sql.functions._
    
      val d1 = spark.sparkContext.parallelize(Seq(
        ("A", "U", "PIERRE", 1),
        ("A", "U", "THOMAS", 2),
        ("A", "U", "MICHAEL", 3),
        ("A", "V", "TOM", 2),
        ("A", "V", "JACK", 3),
        ("A", "W", "MICHEL", 2),
        ("A", "W", "JULIEN", 3)
      )).toDF("ID1", "ID2", "VAL1", "VAL2")
    
    
      d1.groupBy("ID1", "ID2").agg(first(struct("VAL1", "VAL2")).as("val"))
        .select("ID1", "ID2", "val.*")
        .show(false)
    

    更新: 如果你有键和值作为参数,那么你可以使用如下。

    val keys = Seq("ID1", "ID2")
    
    val values = Seq("VAL1", "VAL2")
    
    d1.groupBy(keys.head, keys.tail : _*)
        .agg(first(struct(values.head, values.tail:_*)).as("val"))
        .select( "val.*",keys:_*)
        .show(false)
    

    输出:

    +---+---+------+----+
    |ID1|ID2|VAL1  |VAL2|
    +---+---+------+----+
    |A  |W  |MICHEL|2   |
    |A  |V  |TOM   |2   |
    |A  |U  |PIERRE|1   |
    +---+---+------+----+
    

    我希望这会有所帮助!

    【讨论】:

    • 谢谢!它确实有帮助!但这并不是我想要的,因为我还有另一个困难:即使我可能不知道架构,我也想这样做......我会尝试适应它。
    • 我更新了我的问题,希望这能更准确地引导您达到我的预期!
    猜你喜欢
    • 2016-06-10
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-12-30
    • 1970-01-01
    • 2015-08-11
    • 1970-01-01
    相关资源
    最近更新 更多