【问题标题】:Iterate through rows in DataFrame and transform one to many遍历 DataFrame 中的行并将一对多转换
【发布时间】:2017-10-20 08:20:58
【问题描述】:

作为 scala 中的示例,我有一个列表,并且每个项目都匹配我想要出现两次的条件(可能不是这个用例的最佳选择 - 但想法很重要):

l.flatMap {
  case n if n % 2 == 0 => List(n, n)
  case n => List(n)
}

我想在 Spark 中做类似的事情 - 迭代 DataFrame 中的行,如果一行匹配某个条件,那么我需要复制该行并在副本中进行一些修改。如何做到这一点?

例如,如果我的输入是下表:

| name  | age |
|-------|-----|
| Peter | 50  |
| Paul  | 60  |
| Mary  | 70  |

我想遍历表并针对多个条件测试每一行,并且对于每个匹配的条件,应该使用匹配条件的名称创建一个条目。

例如条件 #1 是“年龄 > 60”,条件 #2 是“name.length

| name  | age |condition|
|-------|-----|---------|
| Paul  | 60  |    2    |
| Mary  | 70  |    1    |
| Mary  | 70  |    2    |

【问题讨论】:

  • 您应该也可以使用flatMap 来做到这一点。你能展示一些实际数据吗?
  • 添加示例以使其更清晰
  • 您想在name.length > 4 处删除行吗,如果age > 60name.length > 4 怎么办?您还需要 condition 列吗?
  • 对于行匹配的每个条件,结果表中都应该有一个条目。如果没有匹配,则没有条目,多个匹配意味着多个条目

标签: scala apache-spark dataframe


【解决方案1】:

你可以filter匹配条件dataframes然后最后union所有这些。

import org.apache.spark.sql.functions._
val condition1DF = df.filter($"age" > 60).withColumn("condition", lit(1))
val condition2DF = df.filter(length($"name") <= 4).withColumn("condition", lit(2))

val finalDF = condition1DF.union(condition2DF)

你应该有你想要的输出

+----+---+---------+
|name|age|condition|
+----+---+---------+
|Mary|70 |1        |
|Paul|60 |2        |
|Mary|70 |2        |
+----+---+---------+

希望回答对你有帮助

【讨论】:

  • 谢谢,这个解决方案给了我最大的灵活性,非常适合多种条件、从配置中读取的条件、多个表等
【解决方案2】:

您还可以使用 UDF 和 explode() 的组合,如下例所示:

// set up example data
case class Pers1 (name:String,age:Int)
val d = Seq(Pers1("Peter",50), Pers1("Paul",60), Pers1("Mary",70))
val df = spark.createDataFrame(d)

// conditions logic - complex as you'd like
// probably should use a Set instead of Sequence but I digress..
val conditions:(String,Int)=>Seq[Int] =  { (name,age) => 
    (if(age > 60) Seq(1) else Seq.empty) ++ 
    (if(name.length <=4) Seq(2) else Seq.empty)  
}
// define UDF for spark
import org.apache.spark.sql.functions.udf
val conditionsUdf = udf(conditions)
// explode() works just like flatmap
val result  = df.withColumn("condition", 
   explode(conditionsUdf(col("name"), col("age"))))
result.show

+----+---+---------+
|name|age|condition|
+----+---+---------+
|Paul| 60|        2|
|Mary| 70|        1|
|Mary| 70|        2|
+----+---+---------+

【讨论】:

    【解决方案3】:

    这是使用rdd.flatMap 将其展平的一种方法:

    import org.apache.spark.sql.types._
    import org.apache.spark.sql.Row
    
    val new_rdd = (df.rdd.flatMap(r => {
        val conditions = Seq((1, r.getAs[Int](1) > 60), (2, r.getAs[String](0).length <= 4))
        conditions.collect{ case (i, c) if c => Row.fromSeq(r.toSeq :+ i) }
    }))
    
    val new_schema = StructType(df.schema :+ StructField("condition", IntegerType, true))
    
    spark.createDataFrame(new_rdd, new_schema).show
    +----+---+---------+
    |name|age|condition|
    +----+---+---------+
    |Paul| 60|        2|
    |Mary| 70|        1|
    |Mary| 70|        2|
    +----+---+---------+
    

    【讨论】:

    • 你为什么用df.rdd来使用flatMap
    • @JacekLaskowski 如果我不使用rdd,我会收到错误消息。 无法找到编码器....
    • @Psidom 试试import spark.implicits._,用它你应该可以在数据帧上使用flatMapspark 这里是SparkSession
    猜你喜欢
    • 2018-12-28
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-12-15
    • 2010-11-10
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多