【问题标题】:Spark convert column filed type which is in Json to multiple rows or nested rowsSpark将Json中的列字段类型转换为多行或嵌套行
【发布时间】:2018-10-03 21:48:51
【问题描述】:

我有下面这样的 json,这只是数据的一部分。所以实际压缩的 json 如果这些数据有很多

 {
    "filed1": "value1",
    "filed2": "value2",
    "data":"{\"info\":[{\"type\":[\"Extra\"],\"value\":9},{\"type\":[\"Free\"],\"value\":8},{\"type\":[\"Actual\"],\"value\":100}]}",
    "code": "0000"
}
{
    "filed1": "value3",
    "filed2": "value4",
    "data":"{\"info\":[{\"type\":[\"Extra\"],\"value\":1001}]}",
    "code": "0001"
}

{
    "filed1": "value5",
    "filed2": "value6",
    "data":"{\"info\":[{\"type\":[\"Actual\"],\"value\":90},{\"type\":[\"Free\"],\"value\":80}]}",
    "code": "0003"
}

当我在 spark 中读取数据列时,数据列被读取为字符串,所以我需要解析并制作如下所示的列,这里每一行都需要转换为多行

filed1   filed2  code  type    Value
value1   value2  0000  Extra   9
value1   value2  0000  Free    8
value1   value2  0000  Actual  100
value3   value4  0001  Extra   1001
value5   value6  0003  Actual  90
value5   value6  0003  Free    80

我在 udfs 下面写了,但我不知道如何为输入的单行创建多行

val getTypeName = udf((strs:String) => {
 // parse json and return types
  })

val getValue = udf((strs:String) => {
 // parse json and return values
  })

val df = spark.read.json("<pathtojson">)
val df1 = df.withColumn("type", getTypeName("data")).withColumn("value", getValue("data"))

但是通过逻辑我只能得到单行,我希望它根据我的数据字段转换两个行数

【问题讨论】:

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


    【解决方案1】:

    您要查找的函数名为explode。基本上,您想编写一个解析 JSON 并输出嵌套行数组的 UDF。然后你在该列上调用explode。 Explode 将采用具有多个值的列并为每个值创建一个新行(复制其他列中的值)。例如:

    case class DataRow(filed1: String, filed2: String, data: String, code: String)
    
    val df = Seq(
        DataRow(
            "value1",
            "value2",
            "{\"info\":[{\"type\":[\"Extra\"],\"value\":9},{\"type\":[\"Free\"],\"value\":8},{\"type\":[\"Actual\"],\"value\":100}]}",
            "0000"
        ),
        DataRow(
            "value3",
            "value4",
            "{\"info\":[{\"type\":[\"Extra\"],\"value\":1001}]}",
            "0001"
        )
    ).toDF 
    
    case class NestedRow(row_type: String, value: Int)
    def processJsonFn(json: String): Seq[NestedRow] = {
        // ... Parse json ...
        val parsed = Seq(NestedRow("Extra", 9), NestedRow("Actual", 100))
    
        parsed
    }
    val processJson = udf(processJsonFn _)
    
    // Convert string json to nested rows
    val df2 = df.withColumn("data", processJson($"data"))
    
    // Explode them
    val df3 = df2.withColumn("data", explode($"data"))
    
    // Flatten structure
    val df4 = df3.select($"filed1", $"filed2", $"data.row_type" as "type", $"data.value" as "value", $"code")
    
    df4.printSchema
    df4.show
    

    输出这个:

    root
     |-- filed1: string (nullable = true)
     |-- filed2: string (nullable = true)
     |-- type: string (nullable = true)
     |-- value: integer (nullable = true)
     |-- code: string (nullable = true)
    
    
    scala> df4.show
    +------+------+------+-----+----+
    |filed1|filed2|  type|value|code|
    +------+------+------+-----+----+
    |value1|value2| Extra|    9|0000|
    |value1|value2|Actual|  100|0000|
    |value3|value4| Extra|    9|0001|
    |value3|value4|Actual|  100|0001|
    +------+------+------+-----+----+
    

    【讨论】:

      猜你喜欢
      • 2019-09-28
      • 2020-02-17
      • 2015-12-30
      • 2018-09-25
      • 2018-12-09
      • 2019-03-28
      • 2019-11-07
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多