【问题标题】:Convert a column in Dataset which have key value pairs into different rows将数据集中具有键值对的列转换为不同的行
【发布时间】:2020-01-29 08:37:59
【问题描述】:

我在一个 dataframe 中有数据,它是从 azure eventthub 获得的。 然后我将此数据转换为 json 对象并将所需的数据存储到数据集中,如下所示。

从 eventthub 获取数据并将其存储到数据框中的代码。

val connectionString = ConnectionStringBuilder(<ENDPOINT URL>)
    .setEventHubName(<EVENTHUB NAME>).build

val currTime = Instant.now
val ehConf = EventHubsConf(connectionString)
    .setConsumerGroup("<CONSUMER GRP>")
    .setStartingPosition(EventPosition
             .fromEnqueuedTime(currTime.minus(Duration.ofMinutes(30))))
    .setEndingPosition(EventPosition.fromEnqueuedTime(currTime))

val reader =  spark.read.format("eventhubs").options(ehConf.toMap).load()

var SIGNALS =  reader
    .select(get_json_object(($"body").cast("string"),"$.NUM").alias("NUM"),
            get_json_object(($"body").cast("string"),"$.SIG1").alias("SIG1"),
            get_json_object(($"body").cast("string"),"$.SIG2").alias("SIG2"),
            get_json_object(($"body").cast("string"),"$.SIG3").alias("SIG3"),
            get_json_object(($"body").cast("string"),"$.SIG4").alias("SIG4")
     )

val SIGNALSFiltered = SIGNALS.filter(col("SIG1").isNotNull &&
    col("SIG2").isNotNull && col("SIG3").isNotNull && col("SIG4").isNotNull)

SIGNALSFiltered处得到的数据如下所示。

+-----------------+--------------------+--------------------+--------------------+--------------------+
|              NUM|                SIG1|                SIG2|                SIG3|                SIG4|
+-----------------+--------------------+--------------------+--------------------+--------------------+
|XXXXX01|[{"TIME":15695605310...|[{"TIME":15695605310...|[{"TIME":15695605310...|[{"TIME":15695605310...|
|XXXXX02|[{"TIME":15695604780...|[{"TIME":15695604780...|[{"TIME":15695604780...|[{"TIME":15695604780...|
|XXXXX03|[{"TIME":15695605310...|[{"TIME":15695605310...|[{"TIME":15695605310...|[{"TIME":15695605310...|
|XXXXX04|[{"TIME":15695605310...|[{"TIME":15695605310...|[{"TIME":15695605310...|[{"TIME":15695605310...|
|XXXXX05|[{"TIME":15695605310...|[{"TIME":15695605310...|[{"TIME":15695605310...|[{"TIME":15695605310...|
|XXXXX06|[{"TIME":15695605340...|[{"TIME":15695605340...|[{"TIME":15695605340...|[{"TIME":15695605340...|
|XXXXX07|[{"TIME":15695605310...|[{"TIME":15695605310...|[{"TIME":15695605310...|[{"TIME":15695605310...|
|XXXXX08|[{"TIME":15695605310...|[{"TIME":15695605310...|[{"TIME":15695605310...|[{"TIME":15695605310...|

如果我们检查单个行的整个数据,它将如下所示。

|XXXXX01|[{"TIME":1569560531000,"VALUE":3.7825},{"TIME":1569560475000,"VALUE":3.7812},{"TIME":1569560483000,"VALUE":1.7812},{"TIME":1569560491000,"VALUE":7.7875}]|
    [{"TIME":1569560537000,"VALUE":3.7825},{"TIME":1569560481000,"VALUE":9.7825},{"TIME":1569560489000,"VALUE":5.7825},{"TIME":1569560497000,"VALUE":34.7825}]|
    [{"TIME":1569560505000,"VALUE":34.7825},{"TIME":1569560513000,"VALUE":9.7825},{"TIME":1569560521000,"VALUE":34.7825},{"TIME":1569560527000,"VALUE":4.7825}]|
    [{"TIME":1569560535000,"VALUE":7.7825},{"TIME":1569560479000,"VALUE":35.7825},{"TIME":1569560487000,"VALUE":3.7825}]

我想将每个信号列中的每个时间值对转换为新行。

有什么方法可以将基础数据集转换如下?列中的每个元素都应转换为新行。

+-----------------+-----------------------------+---------------------------------------+-----------------------------+
|    NUM|    SIG1 TIME| SIG1 VALUE|    SIG2 TIME|   SIG2 VALUE|    SIG3 TIME|   SIG3 VALUE|    SIG4 TIME|  SIG4 VALUE |
+-----------------+-----------------------------+---------------------------------------+-----------------------------+
|XXXXX01|1569560531000|     3.7825|1569560531000|       4.7825|1569560531000|       8.7825|1569560531000|       2.7825|
|XXXXX01|1569560531000|     1.7825|1569560531000|       1.7825|        null |       null  |1569560531000|       2.7825|
|XXXXX01|1569560531000|     3.7825|1569560531000|       4.7825|1569560531000|       8.7825|1569560531000|       7.7825|
|XXXXX02|1569560531000|     7.7825|1569560531000|       4.7825|1569560531000|       8.7825|1569560531000|       2.7825|
|XXXXX02|null         |     null  |1569560531000|       5.7825|1569560531000|       7.7825|1569560531000|       5.7825|
|XXXXX02|1569560531000|     3.7825|1569560531000|       4.7825|1569560531000|       8.7825|1569560531000|       2.7825|
|XXXXX02|1569560531000|     5.7825|1569560531000|       7.7825|1569560531000|       9.7825|1569560531000|       2.7825|

感谢任何线索或帮助!提前致谢。

【问题讨论】:

    标签: dataframe apache-spark apache-spark-sql apache-spark-dataset


    【解决方案1】:
    scala> SIGNALSFiltered.show(false)
    +-------+--------------------------------------------------------------------------------------------------------------+----------------------------------------------------------------------------------------------------------------+-----------------------------------------------------------------------------------------------------------------+-------------------------------------------------------------------------------------+
    |NUM    |SIG1                                                                                                          |SIG2                                                                                                            |SIG3                                                                                                             |SIG4                                                                                 |
    +-------+--------------------------------------------------------------------------------------------------------------+----------------------------------------------------------------------------------------------------------------+-----------------------------------------------------------------------------------------------------------------+-------------------------------------------------------------------------------------+
    |XXXXX01|[{"TIME":11,"VALUE":3.7825},{"TIME":12,"VALUE":3.7812},{"TIME":13,"VALUE":3.7812},{"TIME":14,"VALUE":34.7875}]|[{"TIME":21,"VALUE":3.7825},{"TIME":22,"VALUE":34.7825},{"TIME":23,"VALUE":34.7825},{"TIME":24,"VALUE":34.7825}]|[{"TIME":31,"VALUE":34.7825},{"TIME":32,"VALUE":34.7825},{"TIME":33,"VALUE":34.7825},{"TIME":34,"VALUE":34.7825}]|[{"TIME":41,"VALUE":34.7825},{"TIME":42,"VALUE":34.7825},{"TIME":43,"VALUE":34.7825}]|
    +-------+--------------------------------------------------------------------------------------------------------------+----------------------------------------------------------------------------------------------------------------+-----------------------------------------------------------------------------------------------------------------+-------------------------------------------------------------------------------------+
    
    
    scala>  import scala.collection.mutable.ListBuffer
    scala>  import org.apache.spark.sql.functions.arrays_zip
    scala>  import scala.util.parsing.json._
    
    scala>    def flatTime:UserDefinedFunction = udf((json:String) => {
         |    val pars = JSON.parseFull(json)
         |    var outputList = new ListBuffer[String]()
         |    pars.foreach{ x => 
         |    val y = x.asInstanceOf[List[Any]]
         |    y.foreach{ zz =>
         |    val z =  zz.asInstanceOf[Map[String,Double]]
         |     val tempStr = """[{"TIME" : """ + z("TIME").toString + """ ,"VALUE": """ +  z("VALUE").toString + """}]"""
         |     outputList += tempStr
         |   }
         |   }
         |   outputList.toList
         |   })
    
    scala> SIGNALSFiltered.withColumn("var", explode(arrays_zip(flatTime(col("SIG1")),flatTime(col("SIG2")),flatTime(col("SIG3")),flatTime(col("SIG4"))))).select(col("NUM"), col("var.0").alias("SIG1"),col("var.1").alias("SIG2"),col("var.2").alias("SIG3"),col("var.3").alias("SIG4")).show(false)
    +-------+-----------------------------------+-----------------------------------+-----------------------------------+-----------------------------------+
    |NUM    |SIG1                               |SIG2                               |SIG3                               |SIG4                               |
    +-------+-----------------------------------+-----------------------------------+-----------------------------------+-----------------------------------+
    |XXXXX01|[{"TIME" : 11.0 ,"VALUE": 3.7825}] |[{"TIME" : 21.0 ,"VALUE": 3.7825}] |[{"TIME" : 31.0 ,"VALUE": 34.7825}]|[{"TIME" : 41.0 ,"VALUE": 34.7825}]|
    |XXXXX01|[{"TIME" : 12.0 ,"VALUE": 3.7812}] |[{"TIME" : 22.0 ,"VALUE": 34.7825}]|[{"TIME" : 32.0 ,"VALUE": 34.7825}]|[{"TIME" : 42.0 ,"VALUE": 34.7825}]|
    |XXXXX01|[{"TIME" : 13.0 ,"VALUE": 3.7812}] |[{"TIME" : 23.0 ,"VALUE": 34.7825}]|[{"TIME" : 33.0 ,"VALUE": 34.7825}]|[{"TIME" : 43.0 ,"VALUE": 34.7825}]|
    |XXXXX01|[{"TIME" : 14.0 ,"VALUE": 34.7875}]|[{"TIME" : 24.0 ,"VALUE": 34.7825}]|[{"TIME" : 34.0 ,"VALUE": 34.7825}]|null                               |
    +-------+-----------------------------------+-----------------------------------+-----------------------------------+-----------------------------------+
    

    【讨论】:

      【解决方案2】:

      您可以使用explode 函数来实现。它将为数组中的每个元素生成新行,然后您可以使用点语法访问字段 timevalue(访问结构的字段)。这是第一列的一个简单示例:

      data
      .withColumn("sig1_obj", explode($"SIG1"))
      .withColumn("sig1_time", $"sig1_obj.time")
      .withColumn("sig1_value", $"sig1_obj.value")
      .show()
      
      +--------------------+--------------------+-------------+----------+
      |                SIG1|            sig1_obj|    sig1_time|sig1_value|
      +--------------------+--------------------+-------------+----------+
      |[[1569560531000, ...|[1569560531000, 3...|1569560531000|    3.7825|
      |[[1569560531000, ...|[1569560475000, 3...|1569560475000|    3.7812|
      |[[1569560531000, ...|[1569560483000, 1...|1569560483000|    1.7812|
      |[[1569560531000, ...|[1569560491000, 7...|1569560491000|    7.7875|
      +--------------------+--------------------+-------------+----------+
      

      同样,您也可以处理其他列。

      还请注意,使用此技术将乘以数据,对于第二列,您将获得 n*m 行,其中 n 是 sig1 数组中的元素数,m 是sig2 数组等等。如果您不希望这样,您可以在单独的数据框中分解每一列,然后在某些字段上完全外部连接这些数据框(可能 row_number 为每个 NUM 的行并加入 NUM col 和 row_number)

      编辑:

      由于您的 sig 列中有 StringType,您可以先使用 from_json 函数将此字符串字段转换为结构数组。在您的示例中,可以按如下方式完成:

      import org.apache.spark.sql.types.{StructField, StructType, ArrayType, StringType}
      
      val schema = ArrayType(StructType(Seq(StructField("TIME", StringType), StructField("VALUE", StringType))))
      
      df.withColumn("sig1_arr", from_json($"SIG1", schema))
      
      

      【讨论】:

      • 是的,我也尝试了爆炸功能,但由于信号类型不匹配而失败。信号不是数组或映射类型。 scala&gt; SIGNALSFiltered.printSchema root |-- NUM: string (nullable = true) |-- SIG1: string (nullable = true) |-- SIG2: string (nullable = true) |-- SIG3: string (nullable = true) |-- SIG4: string (nullable = true)@David Vrba
      • @Antony 好的,但是如果字符串是结构数组的形式,则可以使用 from_json 函数和提供的模式进行转换。查看更新的答案
      • 由于我必须转换的列是我上面评论的字符串类型,因此此解决方案不起作用。在应用此解决方案之前,必须将 Signal 列从 String 类型转换为 Struct/Array。 @David Vrba
      • @Antony 抱歉,指定的架构错误,我编辑了答案并更改了架构。它现在对我有用。
      • @Antony 使用正确的模式很重要,否则结果将是空值。
      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2018-11-16
      • 1970-01-01
      • 2019-11-18
      • 2022-11-10
      • 2018-08-07
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多