【问题标题】:Filter based on JSON data which is in a string column in a Spark dataframe基于 Spark 数据帧中字符串列中的 JSON 数据进行过滤
【发布时间】:2021-06-20 23:48:16
【问题描述】:

我有一个以下格式的 Spark 数据框,其中 FamilyDetails 列是一个字符串字段:

root
 |-- FirstName: string (nullable = true)
 |-- LastName: string (nullable = true)
 |-- FamilyDetails: string (nullable = true)

+----------+---------+--------------------------------------------------------------------------------------------------------------------------+
|FirstName |LastName |FamilyDetails                                                                                                             |
+----------+---------+--------------------------------------------------------------------------------------------------------------------------+
|Emma      |Smith    |{                                                                                                                         |
|          |         | "23214598.31601190":{"gender":"F","Name":"Ms Olivia Smith","relationship":"Daughter"},                                   |
|          |         | "23214598.23214598":{"gender":"F","Name":"Ms Emma Smith","relationship":null}                                            |
|          |         |}                                                                                                                         |
|Joe       |Williams |{                                                                                                                         |
|          |         |  "2321463.2321463":{"gender":"M","Name":"Mr Joe Williams","relationship":null},                                          |
|          |         |  "2321463.3841483":{"gender":"F","Name":"Mrs Sophia Williams","relationship":"Wife","IsActive":"N"}                      |
|          |         |}                                                                                                                         |
|Liam      |Jones    |{                                                                                                                         |
|          |         |  "2321464.12379942":{"gender":"F","Name":"Miss Patricia Jones","relationship":"Sister"},                                 |
|          |         |  "2321464.2321464":{"gender":"M","Name":"Mr Liam Jones","relationship":null,"IsActive":"Y"}                              |
|          |         |}                                                                                                                         |
+----------+---------+--------------------------------------------------------------------------------------------------------------------------+

我想做的事:

我正在尝试获取我们有不活跃家庭成员的记录 (IsActive='N')。需要注意的是IsActive是可选字段。

预期输出:

+----------+---------+--------------------------------------------------------------------------------------------------------------------------+
|FirstName |LastName |FamilyDetails                                                                                                             |
+----------+---------+--------------------------------------------------------------------------------------------------------------------------+                                                                                                                   |
|Joe       |Williams |{                                                                                                                         |
|          |         |  "2321463.2321463":{"gender":"M","Name":"Mr Joe Williams","relationship":null},                                          |
|          |         |  "2321463.3841483":{"gender":"F","Name":"Mrs Sophia Williams","relationship":"Wife","IsActive":"N"}                      |
|          |         |}                                                                                                                         |
+----------+---------+--------------------------------------------------------------------------------------------------------------------------+

到目前为止我所做的尝试:

由于不知道完整的架构,我尝试从FamilyDetails 列本身创建架构。

import org.apache.spark.sql.functions._
import spark.implicits._
val json_schema = spark.read.json(myDF.select("FamilyDetails").as[String]).schema
println(json_schema)

这给了我:

StructType(
    StructField(23214598.31601190,
        StructType(
            StructField(gender,StringType,true), 
            StructField(Name,StringType,true), 
            StructField(relationship,StringType,true), 
            StructField(IsActive,StringType,true)
        )
        ,true
    )
)

如何摆脱第一个值 (2321463.2321463) 并仅采用 json 架构中的必填字段?或者有没有更简单的方法来过滤IsActive = 'N' 的记录?

【问题讨论】:

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


    【解决方案1】:

    也许你可以通过简单地找到字符串"IsActive":"N"来避免解析JSON:

    val df2 = df.filter("""FamilyDetails rlike '"IsActive":"N"'""")
    

    为了更严格的解析,可以使用:

    val df2 = df.filter("exists(map_values(from_json(FamilyDetails, 'map<string,map<string,string>>')), x -> x['IsActive'] = 'N')")
    

    【讨论】:

    • 是的,这将是一个更简单的解决方案!但是,我需要对数据进行更多检查。例如:gender 包含 ["M","F"] 以外的值的记录,姓名长度超过 20 个字符的记录,年龄超过 3 位的记录等。
    • 有没有办法从字符串中访问特定的键值对,例如"IsActive"="N""gender"="M"?使用正则表达式什么的?
    • @user2538559 查看编辑后的答案。您可以使用一些 Spark SQL 函数。
    猜你喜欢
    • 2021-12-21
    • 1970-01-01
    • 2022-08-17
    • 1970-01-01
    • 1970-01-01
    • 2021-11-21
    • 2018-06-25
    • 2020-05-04
    • 2021-10-27
    相关资源
    最近更新 更多