【问题标题】:how to update/transform/replace spark df column values using a hashmap如何使用哈希图更新/转换/替换 spark df 列值
【发布时间】:2020-06-25 12:29:34
【问题描述】:

我想使用哈希图替换给定 df 列的值,但我在语法上遇到了困难。 有人可以指出我正确的方向或现有的例子吗?我搜索过但是 找不到能阐明确切主题的东西。

编辑

想象一个如下图所示的数据框:

+-----------+--------+-----------+
|       Noun| Pronoun|  Adjective|
+-----------+--------+-----------+
|      Homer| Simpson|BeerDrinker|
|      Marge| Simpson|  Housewife|
|       Bart| Simpson|        Son|
|       Lisa| Simpson|   Daughter|
|TheSimpsons|Simpsons|     Family|
+-----------+--------+-----------+

我有一个如下所示的键值对映射:

  type ValueMap = scala.collection.mutable.HashMap [String,String]
  var mymap = new ValueMap ()
  mymap += ("Simpson" -> "Surname")

我想做一个操作(我目前还无法弄清楚)并获得如下所示的结果。所以基本上在Pronoun 列中,所有等于Simpson 的列值都已替换为来自映射mymap 的对应值Surname

+-----------+--------+-----------+
|       Noun| Pronoun|  Adjective|
+-----------+--------+-----------+
|      Homer| Surname|BeerDrinker|
|      Marge| Surname|  Housewife|
|       Bart| Surname|        Son|
|       Lisa| Surname|   Daughter|
|TheSimpsons|Simpsons|     Family|
+-----------+--------+-----------+

【问题讨论】:

  • 你能澄清一下替换是什么意思吗?也许提供一个输入数据是什么以及结果应该是什么的例子
  • 嗨@AlexOtt,我已经用一个例子更新了我的程序。我希望这有助于突出我的目标。
  • 我认为这应该回答您的问题:stackoverflow.com/questions/56557771/…(请参阅有关处理不存在密钥的评论)
  • 嗨@Alex:谢谢。是的,它确实解决了问题。不幸的是,我一直在搜索“HashMap”而不是“Map”——只有我这样做了——我才会找到上述帖子。没关系。感谢您向我指出这一点。非常感激!干杯。

标签: scala dataframe apache-spark hashmap transform


【解决方案1】:

用 UDF 试试这个方法,

val myMap = Map("Simpson" -> "Surname")
val df = Seq(("Homer","Simpson","BeerDrinker"),("Marge","Simpson","Housewife"),("Bart","Simpson","Son"),("Lisa","Simpson","Daughter"),("TheSimpsons","Simpsons","Family")).toDF("Noun","Pronoun","Adjective")

df.show(false)

-----------+--------+-----------+
|Noun       |Pronoun |Adjective  |
+-----------+--------+-----------+
|Homer      |Simpson |BeerDrinker|
|Marge      |Simpson |Housewife  |
|Bart       |Simpson |Son        |
|Lisa       |Simpson |Daughter   |
|TheSimpsons|Simpsons|Family     |
+-----------+--------+-----------+

val getVal = udf((x: String) => myMap.getOrElse(x, x))
val resDF = df.withColumn("Pronoun", getVal($"Pronoun"))

resDF.show(false)

+-----------+--------+-----------+
|Noun       |Pronoun |Adjective  |
+-----------+--------+-----------+
|Homer      |Surname |BeerDrinker|
|Marge      |Surname |Housewife  |
|Bart       |Surname |Son        |
|Lisa       |Surname |Daughter   |
|TheSimpsons|Simpsons|Family     |
+-----------+--------+-----------+

如果这有帮助,请告诉我。

更新:

没有UDF,

将地图作为另一列添加到 DF

val df1 = df.withColumn("map", typedLit(myMap))
val df2 = df1.withColumn("Pronoun", when($"map"($"Pronoun").isNotNull, $"map"($"Pronoun")).otherwise($"Pronoun") ).drop("map")
df2.show(false)

+-----------+--------+-----------+
|Noun       |Pronoun |Adjective  |
+-----------+--------+-----------+
|Homer      |Surname |BeerDrinker|
|Marge      |Surname |Housewife  |
|Bart       |Surname |Son        |
|Lisa       |Surname |Daughter   |
|TheSimpsons|Simpsons|Family     |
+-----------+--------+-----------+

另一种简单的方法,而不是添加新列,

val colMap = typedLit(myMap)
val df3 = df.withColumn("Pronoun", when(colMap($"Pronoun").isNotNull, colMap($"Pronoun")).otherwise($"Pronoun") )
df3.show(false)

【讨论】:

  • 嗨,没有UDF就没有办法了吗?不使用 UDF 的原因是在运行 spark 集群时,它不会被催化剂优化器触摸和优化。
  • 更新了 mt 另一种没有 UDF 的方法
  • 嗨@Sathiyan,没有UDF的方法非常完美。我相信这可能是我一直在寻找的。但不幸的是,我无法让它工作。我知道为什么(可能)——因为我不知道“typeLit”以及在这种情况下如何利用它。在此示例之前,我不知道可以将地图转换为列并以所示方式使用它。感谢您花费时间和精力帮助我解决问题。我非常感谢。干杯!
猜你喜欢
  • 1970-01-01
  • 2021-02-12
  • 2013-07-21
  • 2013-02-04
  • 1970-01-01
  • 1970-01-01
  • 2010-12-02
  • 1970-01-01
  • 2011-12-14
相关资源
最近更新 更多