【问题标题】:How to convert map values using Spark built-in functions?如何使用 Spark 内置函数转换地图值?
【发布时间】:2020-03-12 09:49:55
【问题描述】:

我有一个带有架构的简单数据集:

root
 |-- columns: map (nullable = true)
 |    |-- key: string
 |    |-- value: struct (valueContainsNull = false)
 |    |    |-- a: integer (nullable = true)
 |    |    |-- b: long (nullable = true)
 |    |    |-- c: float (nullable = true)
 |    |    |-- d: double (nullable = true)
           |-- ... 
           |-- ...

例子:

+---------------------------------------------------+
|columns                                            |
+---------------------------------------------------+
|[k0 -> [,,,, 2,], k1 -> [,,,, AB,], k2 -> [,,M,,,] |
+---------------------------------------------------+

我想将我的数据集转换为带有架构的新数据集:

root
 |-- columns: map (nullable = true)
 |    |-- key: string
 |    |-- value: string

转换规则:

  1. 结构大小未定义。
  2. 从值结构中获取第一个非空元素(作为字符串)。

输出示例:

+----------------------------+
|columns                     |
+----------------------------+
|[k0 -> 2, k1 -> AB, k2 -> M |
+----------------------------+

这是我的 UDF 解决方案

val my_udf: UserDefinedFunction = udf((m: Map[String, Row]) => m.map { case (k, v) => (k, v.toSeq.find(_ != null).map(_.toString)) })

df.select(my_udf(col("columns")))

是否可以使用 Spark 内置函数重写它?

类似这样的:

df.withColumn("data", expr("transform(fields.items(), (k, v) -> (k, get-1st-not-null-element-from-v)"))

这是另一个尝试(Spark 3.0+):

df.select(map_entries(col("fields")).as("array"))
.select(
  expr(
    "transform(array, (e, _) -> " +
      "struct(cast(e.key as string), coalesce(e.value.a, e.value.b, e.value.c, e.value.d, ...)))"
    ).as("entries")
  )
.select(map_from_entries(col("entries")))

【问题讨论】:

  • 如果您使用的是 Spark 2.4+,那么您是否考虑过一个高阶函数 - transform()selectExpr() 内部传递?
  • 我认为 Spark 中没有可用的 map (k, v) 迭代。 UDF 似乎是唯一的方法
  • @undying_odyssey 是的,我在帖子中添加了“转换”示例。但是它使用 Spark 3.0+ 功能。我宁愿继续使用 Spark 2.4.x。

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


【解决方案1】:

IIUC,你可以试试transform + aggregate(假设列名是col1):

df.selectExpr("""
  aggregate(
    transform(map_keys(col1), x -> map(x, coalesce(col1[x].a,col1[x].b,col1[x].c,col1[x].d))), 
    /* zero_value: use an empty map() */
    map(), 
    /* merge: do map_concat() */
    (acc,y) -> map_concat(acc, y)
  )  as col1
""").show()

地点:

  • 使用transform()函数从map_keys遍历数组,将x的每一项转换成以x为key的map,并将值设置为第一个使用 coalesce(col1[x].a,col1[x].b,col1[x].c,col1[x].d) 的 StructType 字段中的非空值。这将产生一个地图数组。

  • 使用 aggregate() 函数将上述地图数组合并为一个 MapType 列。

对于 spark 3.0+,使用 transform_values:

df.selectExpr("transform_values(col1, (k,v) -> coalesce(v.a, v.b, v.c, v.d)) as col1").show()

【讨论】:

  • OP 指定结构大小未定义。因此,您必须查看数据框的架构以查找结构字段名称。此外,如果您使用的是 Spark 3.0,请将 DSL 用于高阶函数:比生成 SQL 的代码更容易。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2015-11-11
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多