【问题标题】:spark dataframes : reading json having duplicate column names but different datatypesspark dataframes:读取具有重复列名但数据类型不同的 json
【发布时间】:2020-10-14 22:49:57
【问题描述】:

我有下面这样的 json 数据,其中版本字段是差异化因素 -

file_1 = {"version": 1, "stats": {"hits":20}}

file_2 = {"version": 2, "stats": [{"hour":1,"hits":10},{"hour":2,"hits":12}]}

在新格式中,stats 列现在是Arraytype(StructType)。

之前只需要file_1,所以我使用

spark.read.schema(schema_def_v1).json(path)

现在我需要同时读取这两种类型的多个 json 文件。我无法在 schema_def 中将 stats 定义为字符串,因为这会影响 corruptrecord 功能(用于 stats 列),该功能检查所有字段的格式错误的 json 和架构合规性。

1 个只读文件中需要的示例 df 输出 -

version | hour | hits
1       | null | 20
2       | 1    | 10
2       | 2    | 12

我尝试使用mergeSchema 选项读取,但这会使统计字段成为字符串类型。

另外,我尝试通过过滤版本字段并应用spark.read.schema(schema_def_v1).json(df_v1.toJSON) 来制作两个数据帧。这里的 stats 列也变成了String 类型。

我在想,如果在阅读时,我可以根据数据类型将 df 列标题解析为 stats_v1 和 stats_v2 可以解决问题。请提供任何可能的解决方案。

【问题讨论】:

    标签: json apache-spark apache-spark-sql jackson jsonschema


    【解决方案1】:

    UDF检查字符串或数组,如果是字符串则将字符串转换为数组。

    import org.apache.spark.sql.functions.udf
    import org.json4s.{DefaultFormats, JObject}
    import org.json4s.jackson.JsonMethods.parse
    import org.json4s.jackson.Serialization.write
    import scala.util.{Failure, Success, Try}
    
    object Parse {
        implicit val formats = DefaultFormats
        def toArray(data:String) = {
          val json_data = (parse(data))
          if(json_data.isInstanceOf[JObject]) write(List(json_data)) else data
        }
    }
    
    val toJsonArray = udf(Parse.toArray _)
    
    
    scala> "ls -ltr /tmp/data".!
    total 16
    -rw-r--r--  1 srinivas  root  37 Jun 26 17:49 file_1.json
    -rw-r--r--  1 srinivas  root  69 Jun 26 17:49 file_2.json
    res4: Int = 0
    
    scala> val df = spark.read.json("/tmp/data").select("stats","version")
    df: org.apache.spark.sql.DataFrame = [stats: string, version: bigint]
    
    scala> df.printSchema
    root
     |-- stats: string (nullable = true)
     |-- version: long (nullable = true)
    
    
    scala> df.show(false)
    +-------+-------------------------------------------+
    |version|stats                                      |
    +-------+-------------------------------------------+
    |1      |{"hits":20}                                |
    |2      |[{"hour":1,"hits":10},{"hour":2,"hits":12}]|
    +-------+-------------------------------------------+
    

    输出

    scala> 
    
    import org.apache.spark.sql.types._
    val schema = ArrayType(MapType(StringType,IntegerType))
    
    df
    .withColumn("json_stats",explode(from_json(toJsonArray($"stats"),schema)))
    .select(
        $"version",
        $"stats",
        $"json_stats".getItem("hour").as("hour"),
        $"json_stats".getItem("hits").as("hits")
    ).show(false)
    
    +-------+-------------------------------------------+----+----+
    |version|stats                                      |hour|hits|
    +-------+-------------------------------------------+----+----+
    |1      |{"hits":20}                                |null|20  |
    |2      |[{"hour":1,"hits":10},{"hour":2,"hits":12}]|1   |10  |
    |2      |[{"hour":1,"hits":10},{"hour":2,"hits":12}]|2   |12  |
    +-------+-------------------------------------------+----+----+
    
    

    没有 UDF

    scala> val schema = ArrayType(MapType(StringType,IntegerType))
    
    scala> val expr = when(!$"stats".contains("[{"),concat(lit("["),$"stats",lit("]"))).otherwise($"stats")
    
    df
    .withColumn("stats",expr)
    .withColumn("stats",explode(from_json($"stats",schema)))
    .select(
        $"version",
        $"stats",
        $"stats".getItem("hour").as("hour"),
        $"stats".getItem("hits").as("hits")
    )
    .show(false)
    
    +-------+-----------------------+----+----+
    |version|stats                  |hour|hits|
    +-------+-----------------------+----+----+
    |1      |[hits -> 20]           |null|20  |
    |2      |[hour -> 1, hits -> 10]|1   |10  |
    |2      |[hour -> 2, hits -> 12]|2   |12  |
    +-------+-----------------------+----+----+
    

    【讨论】:

    • val schema_def = StructType.fromDDL(version int, stats struct OR array> OR string , _corrupt_record 字符串)。我想阅读这样的架构定义,并为格式错误的 json 和架构不匹配收集损坏的记录列。以字符串形式读取统计信息并没有帮助。
    • 你能解释一下为什么阅读字符串没有帮助吗? &上面的例子我读为字符串然后解析成数组。
    • file_1 = {"version": 1, "stats": {"hits":"abcd"}} 。考虑到 hits 字段中这种类型的错误 json 数据,我具有捕获损坏记录的功能在我目前的解决方案中,我相信在这个混合数据类型的新问题中它不会与字符串类型一起使用。您的解决方案有效,但我不想失去这个附加功能。你可以在这里参考(A)部分blog.knoldus.com/apache-spark-handle-corrupt-bad-records
    • 加载时不要应用架构.. 只需加载所有文件,火花就会找出架构。加载后,您可以转换所需的列。在这种情况下,stats 具有不同的数据类型,因此 spark 将添加超类型,即字符串。
    • 你在用python吗??如果是,写udf将字符串转换为字符串数组。然后应用架构它会工作。
    【解决方案2】:

    先读取第二个文件,爆破stats,使用schema读取第一个文件。

    
    from pyspark.sql import SparkSession
    from pyspark.sql.functions import  explode
    
    spark = SparkSession.builder.getOrCreate()
    sc = spark.sparkContext
    
    file_1 = {"version": 1, "stats": {"hits": 20}}
    
    file_2 = {"version": 2, "stats": [{"hour": 1, "hits": 10}, {"hour": 2, "hits": 12}]}
    
    df1 = spark.read.json(sc.parallelize([file_2])).withColumn('stats', explode('stats'))
    schema = df1.schema
    
    spark.read.schema(schema).json(sc.parallelize([file_1])).printSchema()
    
    output >> root
     |-- stats: struct (nullable = true)
     |    |-- hits: long (nullable = true)
     |    |-- hour: long (nullable = true)
     |-- version: long (nullable = true)
    
    

    【讨论】:

    • 谢谢,但我在许多此类文件中没有任何标识。他们都是混合的
    • 另外我认为如果你在代码中显示()df,file_2的stats部分将显示为null,因为使用的模式没有ArrayType
    • 我的错,我会改的
    【解决方案3】:

    IIUC,您可以使用spark.read.text 读取JSON 文件,然后使用json_tuple、from_json 解析value。注意stats 字段我们使用coalesce 来解析基于两个或多个模式的字段。 (如果每个文件包含跨多行的单个 JSON 文档,则添加 wholetext=True 作为 spark.read.text 的参数)

    from pyspark.sql.functions import json_tuple, coalesce, from_json, array
    
    df = spark.read.text("/path/to/all/jsons/")
    
    schema_1 = "array<struct<hour:int,hits:int>>"
    schema_2 = "struct<hour:int,hits:int>"
    
    df.select(json_tuple('value', 'version', 'stats').alias('version', 'stats')) \
        .withColumn('status', coalesce(from_json('stats', schema_1), array(from_json('stats', schema_2)))) \
        .selectExpr('version', 'inline_outer(status)') \
        .show()
    +-------+----+----+
    |version|hour|hits|
    +-------+----+----+
    |      2|   1|  10|
    |      2|   2|  12|
    |      1|null|  20|
    +-------+----+----+
    

    【讨论】:

    • 我正在使用 scala,我刚刚发现我们可以在不合并的情况下解析混合类型统计列,只需使用 from_json 函数中的 ArrayType schema_def
    • 一个重要的需求是在 stats 列中捕获错误的数据类型记录,但我读到带有第三个参数的 from_json 函数 - options Map 在 spark 3.0 之前不支持 PERMISSIVE 模式或 columnNameOfCorruptRecord跨度>
    • 是的,你可以在 from_json 函数中使用array&lt;map&lt;string,string&gt;&gt;,就可以了。但是您确实使用 spark.read.text 读取文件然后解析该字段? @Abhishek
    • 等等,我想还是会错过map&lt;string,string&gt;,类似于结构数组。那么为什么只使用 coalesce 而不是依赖其他选项呢。
    • 我做了 spark.read.json。它给了我作为 StringType 的统计信息,我在 from_json 中应用了 schema_1 来一次解析 stats 字段中的两种类型
    猜你喜欢
    • 1970-01-01
    • 2021-09-20
    • 2012-08-18
    • 2017-09-12
    • 2017-06-28
    • 2018-12-01
    • 2016-01-25
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多