【问题标题】:Working with nested data in Spark SQL / High order functions在 Spark SQL / 高阶函数中处理嵌套数据
【发布时间】:2018-08-15 03:58:24
【问题描述】:

我最近开始使用 Spark SQL (2.1),我正在处理嵌套数据。

这是我的架构:

 root
 |-- a: string (nullable = true)
 |-- b: map (nullable = true)
 |    |-- bb: string
 |    |-- bbb: string (valueContainsNull = true)
 |-- c: array (nullable = true)
 |    |-- element: struct (containsNull = true)
 |    |    |-- cc: map (nullable = true)
 |    |    |    |-- cca: string
 |    |    |    |-- ccb: struct (valueContainsNull = true)
 |    |    |    |    |-- member0: string (nullable = true)
 |    |    |    |    |-- member1: long (nullable = true)
 |    |    |-- ccc: map (nullable = true)
 |    |    |    |-- ccca: string
 |    |    |    |-- cccb: string (valueContainsNull = true)
 |    |    |-- cccc: map (nullable = true)
 |    |    |    |-- cccca: string
 |    |    |    |-- ccccb: string (valueContainsNull = true)

我正在尝试按如下方式过滤我的数据:保留 c.ccc.key == 'data' 的所有行

我发现非常相关的函数存在于 databricks 文档中。但我想知道databricks笔记本之外是否有类似的东西?

https://docs.databricks.com/spark/latest/spark-sql/higher-order-functions-lambda-functions.html#exists-array-t-function-t-v-boolean-boolean

我愿意使用 sql 或以编程方式执行,只是不确定数据帧如何不是类型对象。

阅读此电子邮件线程http://apache-spark-developers-list.1001551.n3.nabble.com/Will-higher-order-functions-in-spark-SQL-be-pushed-upstream-td21703.html 似乎数据块中的高阶函数将很快可供所有人使用。但我想知道是否有任何人可以分享的中间解决方案?

【问题讨论】:

  • 顺便说一句,使用“exists”函数它看起来像这样“SELECT *, EXISTS(c, item -> item.ccc.key == 'data') FROM mytable”
  • 你能分享你的尝试吗?并分享一些样本数据进行测试。
  • 一行应该是这样的:name | Map(k->v) | [(Map(kk->vv), Map(kkk->vvv, key->data), Map(kkk->vvv))]
  • 从 databricks 中查看笔记本示例(数据也可以在此处获得)-docs.databricks.com/spark/latest/spark-sql/… 我正在寻找可以完全充当“存在”功能的东西
  • 基本上对于我表中的每一行,我想遍历数组(字段 c),如果数组中的一个元素遵循某个条件,那么我想保留/标记这个行。

标签: apache-spark apache-spark-sql


【解决方案1】:

如果你的dataframeschema

root
 |-- a: string (nullable = true)
 |-- b: map (nullable = true)
 |    |-- key: string
 |    |-- value: string (valueContainsNull = true)
 |-- c: array (nullable = true)
 |    |-- element: struct (containsNull = true)
 |    |    |-- cc: map (nullable = true)
 |    |    |    |-- key: string
 |    |    |    |-- value: struct (valueContainsNull = true)
 |    |    |    |    |-- member0: string (nullable = true)
 |    |    |    |    |-- member1: string (nullable = true)
 |    |    |-- ccc: map (nullable = true)
 |    |    |    |-- key: string
 |    |    |    |-- value: string (valueContainsNull = true)
 |    |    |-- cccc: map (nullable = true)
 |    |    |    |-- key: string
 |    |    |    |-- value: string (valueContainsNull = true)

然后你可以写一个udf函数如下

import org.apache.spark.sql.functions._
def filterUdf = udf((column: Seq[Row])=> column.map(x => x(1).asInstanceOf[Map[String, String]].keySet.contains("data")).contains(true))

这将扫描每一行c 列是否存在data 字符串,您可以将filter 函数中的udf 函数用作

df.filter(filterUdf(col("c")))

所以最后你应该在c.ccc.key 中只有带有data 的行

【讨论】:

  • @Maayan 我想你也可以投票 :) 如果帖子有帮助 :)
【解决方案2】:

谢谢! 一些改进:

spark.sqlContext.udf.register("contains_key", (field: Seq[Row], key: String, value: String) => field.exists(item => item.getAs[Map[String, String]]("ccc").get(key).getOrElse("").equals(value)))

然后就可以用 spark sql 访问了:

spark.sql("select contains_key(c,"key","data") from mytable")

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-04-26
    • 1970-01-01
    • 1970-01-01
    • 2019-12-27
    相关资源
    最近更新 更多