【问题标题】:Spark - How to apply rules defined in a dataframe to another dataframeSpark - 如何将数据框中定义的规则应用于另一个数据框
【发布时间】:2017-12-16 16:49:43
【问题描述】:

我正在尝试使用 Spark 2 解决此类问题,但找不到解决方案。

我有一个数据框 A

+----+-------+------+
|id  |COUNTRY| MONTH|
+----+-------+------+
|  1 |    US |    1 |
|  2 |    FR |    1 |
|  4 |    DE |    1 |
|  5 |    DE |    2 |
|  3 |    DE |    3 |
+----+-------+------+

还有一个数据框B

+-------+------+------+
|COLUMN |VALUE | PRIO |
+-------+------+------+
|COUNTRY|   US |    5 |
|COUNTRY|   FR |   15 |
|MONTH  |   3  |    2 |
+-------+------+------+

我们的想法是将数据帧 B 的“规则”应用于数据帧 A 以获得此结果:

数据框 A'

+----+-------+------+------+
|id  |COUNTRY| MONTH| PRIO |
+----+-------+------+------+
|  1 |    US |    1 |    5 |
|  2 |    FR |    1 |   15 |
|  4 |    DE |    1 |   20 |
|  5 |    DE |    2 |   20 |
|  3 |    DE |    3 |    2 |
+----+-------+------+------+

我试过类似的东西:

dfB.collect.foreach( r =>
    var dfAp = dfA.where(r.getAs("COLUMN") == r.getAs("VALUE"))
    dfAp.withColumn("PRIO", lit(r.getAs("PRIO")))
)

但我确定这不是正确的方式。

在 Spark 中解决这个问题的策略是什么?

【问题讨论】:

    标签: scala apache-spark apache-spark-sql


    【解决方案1】:

    假设规则集相当小(可能关注数据的大小和生成的表达式的大小,在最坏的情况下,can crash the planner)最简单的解决方案是使用本地收集和将其映射到 SQL 表达式:

    import org.apache.spark.sql.functions.{coalesce, col, lit, when}
    
    val df = Seq(
      (1, "US", "1"), (2, "FR", "1"), (4, "DE", "1"),
      (5, "DE", "2"), (3, "DE", "3")
    ).toDF("id", "COUNTRY", "MONTH")
    
    val rules = Seq(
      ("COUNTRY", "US", 5), ("COUNTRY", "FR", 15), ("MONTH", "3", 2)
    ).toDF("COLUMN", "VALUE", "PRIO")
    
    
    val prio = coalesce(rules.as[(String, String, Int)].collect.map {
      case (c, v, p) => when(col(c) === v, p)
    } :+ lit(20): _*)
    
    df.withColumn("PRIO", prio)
    
    +---+-------+-----+----+
    | id|COUNTRY|MONTH|PRIO|
    +---+-------+-----+----+
    |  1|     US|    1|   5|
    |  2|     FR|    1|  15|
    |  4|     DE|    1|  20|
    |  5|     DE|    2|  20|
    |  3|     DE|    3|   2|
    +---+-------+-----+----+
    

    您可以将coalesce 替换为leastgreatest 以分别应用最小或最大匹配值。

    使用更多的规则,您可以:

    • melt data 转换为长格式。

      val dfLong = df.melt(Seq("id"), df.columns.tail, "COLUMN", "VALUE")
      
    • join 按列和值。

    • 使用适当的聚合函数(例如 min)通过 id 聚合 PRIOR

      val priorities = dfLong.join(rules, Seq("COLUMN", "VALUE"))
        .groupBy("id")
        .agg(min("PRIO").alias("PRIO"))
      
    • 通过id将输出与df外连接。

      df.join(priorities, Seq("id"), "leftouter").na.fill(20)
      
      +---+-------+-----+----+   
      | id|COUNTRY|MONTH|PRIO|
      +---+-------+-----+----+
      |  1|     US|    1|   5|
      |  2|     FR|    1|  15|
      |  4|     DE|    1|  20|
      |  5|     DE|    2|  20|
      |  3|     DE|    3|   2|
      +---+-------+-----+----+
      

    【讨论】:

      【解决方案2】:

      让我们假设 dataframeB 的规则是有限的

      我为下表创建了数据框“df”

      +---+-------+------+
      | id|COUNTRY|MONTH|
      +---+-------+------+
      |  1|     US|     1|
      |  2|     FR|     1|
      |  4|     DE|     1|
      |  5|     DE|     2|
      |  3|     DE|     3|
      +---+-------+------+
      

      通过使用UDF

      val code = udf{(x:String,y:Int)=>if(x=="US") "5" else if (x=="FR") "15" else if (y==3) "2"  else "20"}
      
      df.withColumn("PRIO",code($"COUNTRY",$"MONTH")).show()
      

      输出

      +---+-------+------+----+
      | id|COUNTRY|MONTH|PRIO|
      +---+-------+------+----+
      |  1|     US|     1|   5|
      |  2|     FR|     1|  15|
      |  4|     DE|     1|  20|
      |  5|     DE|     2|  20|
      |  3|     DE|     3|   2|
      +---+-------+------+----+
      

      【讨论】:

        猜你喜欢
        • 2017-01-28
        • 2018-06-04
        • 2015-04-02
        • 2020-02-12
        • 2015-11-14
        • 2020-08-13
        • 2019-02-21
        • 1970-01-01
        • 2022-11-02
        相关资源
        最近更新 更多