【问题标题】:Doing multiple column value look up after joining with lookup dataset加入查找数据集后进行多列值查找
【发布时间】:2020-11-18 02:08:57
【问题描述】:

我正在使用 spark-sql-2.4.1v 如何根据列的值进行各种连接 我需要为给定值列获取map_val 列的多个查找值,如下所示。

样本数据:

val data = List(
  ("20", "score", "school", "2018-03-31", 14 , 12),
  ("21", "score", "school", "2018-03-31", 13 , 13),
  ("22", "rate", "school", "2018-03-31", 11 , 14),
  ("21", "rate", "school", "2018-03-31", 13 , 12)
 )
val df = data.toDF("id", "code", "entity", "date", "value1", "value2")

df.show

+---+-----+------+----------+------+------+
| id| code|entity|      date|value1|value2|
+---+-----+------+----------+------+------+
| 20|score|school|2018-03-31|    14|    12|
| 21|score|school|2018-03-31|    13|    13|
| 22| rate|school|2018-03-31|    11|    14|
| 21| rate|school|2018-03-31|    13|    12|
+---+-----+------+----------+------+------+

查找数据集rateDs

val rateDs = List(
  ("21","2018-01-31","2018-06-31", 12 ,"C"),
  ("21","2018-01-31","2018-06-31", 13 ,"D")
).toDF("id","start_date","end_date", "map_code","map_val")

rateDs.show

+---+----------+----------+--------+-------+
| id|start_date|  end_date|map_code|map_val|
+---+----------+----------+--------+-------+
| 21|2018-01-31|2018-06-31|      12|      C|
| 21|2018-01-31|2018-06-31|      13|      D|
+---+----------+----------+--------+-------+

加入基于start_dateend_datemap_val 列的查找表:

 val  resultDs = df.filter(col("code").equalTo(lit("rate"))).join(rateDs , 
            (
                   df.col("date").between(rateDs.col("start_date"), rateDs.col("end_date"))
                   .and(rateDs.col("id").equalTo(df.col("id"))) 
                   //.and(rateDs.col("mapping_value").equalTo(df.col("mean"))) 
            )
            , "left"
            )
            //.drop("start_date")
            //.drop("end_date")



resultDs.show



+---+----+------+----------+------+------+----+----------+----------+--------+-------+
| id|code|entity|      date|value1|value2|  id|start_date|  end_date|map_code|map_val|
+---+----+------+----------+------+------+----+----------+----------+--------+-------+
| 21|rate|school|2018-03-31|    13|    12|  21|2018-01-31|2018-06-31|      13|      D|
| 21|rate|school|2018-03-31|    13|    12|  21|2018-01-31|2018-06-31|      12|      C|
+---+----+------+----------+------+------+----+----------+----------+--------+-------+

预期的输出应该是:

+---+----+------+----------+------+------+----+----------+----------+--------+-------+
| id|code|entity|      date|value1|value2|  id|start_date|  end_date|map_code|map_val|
+---+----+------+----------+------+------+----+----------+----------+--------+-------+
| 21|rate|school|2018-03-31|    D |    C |  21|2018-01-31|2018-06-31|      13|      D|
| 21|rate|school|2018-03-31|    D |    C |  21|2018-01-31|2018-06-31|      12|      C|
+---+----+------+----------+------+------+----+----------+----------+--------+-------+

如果需要更多详细信息,请告诉我。

【问题讨论】:

  • 不确定我是否理解实际问题
  • 为什么输出包含 4 行的 code=rate 而不是 2?
  • 如果您只对code == rate 的行感兴趣,那么您可以在join 之前过滤掉这些行。
  • 加入按预期工作。如果您需要按预期输出,请在过滤器中添加更多条件。正如@Shaido-ReinstateMonica 所说,在代码列上添加条件将解决您的问题
  • @Shaido - Reinstate Monica 我更正了问题,请检查一次,整个问题是如何查找列,即“value1”、“value2”列,其中包含“map_code”,我需要替换与各自的“map_val”列..

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


【解决方案1】:

试试这个-

Create lookup map before join per id and use the same to replace

 val newRateDS = rateDs.withColumn("lookUpMap",
      map_from_entries(collect_list(struct(col("map_code"), col("map_val"))).over(Window.partitionBy("id")))
    )

    newRateDS.show(false)
    /**
      * +---+----------+----------+--------+-------+------------------+
      * |id |start_date|end_date  |map_code|map_val|lookUpMap         |
      * +---+----------+----------+--------+-------+------------------+
      * |21 |2018-01-31|2018-06-31|12      |C      |[12 -> C, 13 -> D]|
      * |21 |2018-01-31|2018-06-31|13      |D      |[12 -> C, 13 -> D]|
      * +---+----------+----------+--------+-------+------------------+
      */

    val  resultDs = df.filter(col("code").equalTo(lit("rate"))).join(broadcast(newRateDS) ,
      rateDs("id") === df("id") && df("date").between(rateDs("start_date"), rateDs("end_date"))
        //.and(rateDs.col("mapping_value").equalTo(df.col("mean")))
      , "left"
    )

    resultDs.withColumn("value1", expr("coalesce(lookUpMap[value1], value1)"))
      .withColumn("value2", expr("coalesce(lookUpMap[value2], value2)"))
      .show(false)

    /**
      * +---+----+------+----------+------+------+----+----------+----------+--------+-------+------------------+
      * |id |code|entity|date      |value1|value2|id  |start_date|end_date  |map_code|map_val|lookUpMap         |
      * +---+----+------+----------+------+------+----+----------+----------+--------+-------+------------------+
      * |22 |rate|school|2018-03-31|11    |14    |null|null      |null      |null    |null   |null              |
      * |21 |rate|school|2018-03-31|D     |C     |21  |2018-01-31|2018-06-31|13      |D      |[12 -> C, 13 -> D]|
      * |21 |rate|school|2018-03-31|D     |C     |21  |2018-01-31|2018-06-31|12      |C      |[12 -> C, 13 -> D]|
      * +---+----+------+----------+------+------+----+----------+----------+--------+-------+------------------+
      */

【讨论】:

  • 非常感谢,让我看看,coalesce 在这里做什么?
  • 如果没有查找就不会替换
  • 哦,谢谢,如果不查找,要保持为空吗?
  • 为什么这个 .over(Window.partitionBy("id"))) ...over "id" ??
  • 太多查询无法回答 :( 我认为,这已经解决了描述中的任何问题。如何查找。请将此作为参考。您可能想更改如果没有可用的查找值时需要 null,则可以删除 coalesce 之类的东西
猜你喜欢
  • 2012-09-19
  • 2015-09-29
  • 1970-01-01
  • 2021-06-22
  • 2014-10-21
  • 1970-01-01
  • 2018-06-21
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多