【问题标题】:Spark union fails with nested JSON dataframeSpark 联合因嵌套 JSON 数据帧而失败
【发布时间】:2017-03-01 11:24:34
【问题描述】:

我有以下两个 JSON 文件:

{
    "name" : "Agent1",
    "age" : "32",
    "details" : [{
            "d1" : 1,
            "d2" : 2
        }
    ]
}

{
    "name" : "Agent2",
    "age" : "42",
    "details" : []
}

我用火花读过它们:

val jsonDf1 = spark.read.json(pathToJson1)
val jsonDf2 = spark.read.json(pathToJson2)

使用以下架构创建两个数据框:

root
 |-- age: string (nullable = true)
 |-- details: array (nullable = true)
 |    |-- element: struct (containsNull = true)
 |    |    |-- d1: long (nullable = true)
 |    |    |-- d2: long (nullable = true)
 |-- name: string (nullable = true)

root
|-- age: string (nullable = true)
|-- details: array (nullable = true)
|    |-- element: string (containsNull = true)
|-- name: string (nullable = true)

当我尝试对这两个数据框执行联合时,我收到此错误:

jsonDf1.union(jsonDf2)


org.apache.spark.sql.AnalysisException: unresolved operator 'Union;;
'Union
:- LogicalRDD [age#0, details#1, name#2]
+- LogicalRDD [age#7, details#8, name#9]

我该如何解决这个问题?我有时会在 Spark 作业将加载的 JSON 文件中得到空数组,但它仍然必须统一它们,这应该不是问题,因为 Json 文件的架构是相同的。

【问题讨论】:

    标签: scala apache-spark union spark-dataframe


    【解决方案1】:

    如果您尝试合并 2 个数据框,您将得到:

    error:org.apache.spark.sql.AnalysisException: Union can only be performed on tables with the compatible column types. ArrayType(StringType,true) <> ArrayType(StructType(StructField(d1,StringType,true), StructField(d2,StringType,true)),true) at the second column of the second table

    Json 文件同时到达

    为了解决这个问题,如果你能同时读取JSON,我会建议:

    val jsonDf1 = spark.read.json("json1.json", "json2.json")

    这将给出这个架构:

    jsonDf1.printSchema
     |-- age: string (nullable = true)
     |-- details: array (nullable = true)
     |    |-- element: struct (containsNull = true)
     |    |    |-- d1: long (nullable = true)
     |    |    |-- d2: long (nullable = true)
     |-- name: string (nullable = true)
    

    数据输出

    jsonDf1.show(10,truncate = false)
    +---+-------+------+
    |age|details|name  |
    +---+-------+------+
    |32 |[[1,2]]|Agent1|
    |42 |null   |Agent2|
    +---+-------+------+
    

    Json 文件到达的时间不同

    如果您的 json 在不同时间到达,作为默认解决方案,我建议您读取具有完整数组的模板 JSON 对象,这将使您的数据框具有可能对任何联合有效的空数组。然后,您将在输出结果之前使用过滤器删除此假 JSON:

    val df = spark.read.json("jsonWithMaybeAnEmptyArray.json", 
    "TemplateFakeJsonWithAFullArray.json")
    
    df.filter($"name" !== "FakeAgent").show(1)
    

    请注意:已开通 Jira 卡以提高合并 SQL 数据类型的能力:https://issues.apache.org/jira/browse/SPARK-19536,这种操作在下一个 Spark 版本中应该可以实现。

    【讨论】:

    • 谢谢,这确实解决了这个问题,问题是我得到了多个 json 文件(然后转换为 ORC),所以我不能使用这个解决方案。读取函数无法接收文件路径列表或一个连接的路径字符串。所以我必须在 DF 之间进行联合,然后我再次面临同样的问题:\
    • 即使您在路径中使用通配符,像这样:spark.read.json("JsonPathWithWildcards*.json").show(10, truncate = false)?
    【解决方案2】:

    polomarcus 的回答让我想到了这个解决方案: 我无法一次读取所有文件,因为我得到了一个文件列表作为输入,而 spark 没有接收路径列表的 API,但显然使用 Scala 可以做到这一点:

    val files = List("path1", "path2", "path3")
    val dataframe = spark.read.json(files: _*)
    

    这样我得到了一个包含所有三个文件的数据框。

    【讨论】:

      猜你喜欢
      • 2019-04-11
      • 1970-01-01
      • 2019-09-19
      • 2013-02-19
      • 2017-08-13
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多