【问题标题】:Way to concatenate Array of structs连接结构数组的方法
【发布时间】:2020-09-20 02:43:01
【问题描述】:

我有一列包含结构数组。它看起来像这样:

 |-- Network: array (nullable = true)
 |    |-- element: struct (containsNull = true)
 |    |    |-- Code: string (nullable = true)
 |    |    |-- Signal: string (nullable = true)

这只是一个小样本,结构内的列比这多得多。有没有办法为每一行获取列中的数组,将它们连接起来并将它们变成一个字符串?例如,我们可以有这样的东西:

[["example", 2], ["example2", 3]]

有没有办法做成:

"example2example3"?

【问题讨论】:

  • 您使用的是哪个 Spark 版本?

标签: apache-spark


【解决方案1】:

假设有一个数据框df 具有以下架构:

df.printSchema

df 带有样本数据:

df.show(false)

需要先分解Network数组,选择结构体元素Code和signal。

var myDf = df.select(explode($"Network").as("Network"))

然后,您需要使用 concat() 函数连接两列,然后将输出传递给 collect_list() 函数,该函数会将所有行聚合为数组类型的一行

myDf = myDf.select(collect_list(concat($"Network.code",$"Network.signal")).as("data"))

最后,您需要连接成所需的格式,这可以使用 concat_ws() 函数完成,该函数接受两个参数,第一个是要放置在两个字符串之间的分隔符,第二个参数是带有 array 的列type 这是我们上一步的输出。根据您的用例,我们不需要在两个连接字符串之间放置任何分隔符,因此我们将分隔符参数保留为空引号。

myDf = myDf.select(concat_ws("",$"data").as("data"))

以上所有步骤都可以在一行中完成

myDf= myDf.select(explode($"Network").as("Network")).select(concat_ws("",collect_list(concat($"Network.code",$"Network.signal"))).as("data")).show(false)

如果您希望将输出直接转换为字符串变量,请使用:

val myStr = myDf.first.get(0).toString
print(myStr)

【讨论】:

    【解决方案2】:

    有一个名为 spark-hats 的库(Githubsmall article)在这些情况下您可能会发现它非常有用。

    使用它,您可以轻松地映射数组并在元素旁边输出串联,如果您提供完全限定的名称,甚至可以在其他地方输出。

    设置

    import org.apache.spark.sql.functions._
    import za.co.absa.spark.hats.Extensions._
    
    scala> df.printSchema
    root
     |-- info: struct (nullable = true)
     |    |-- drivers: struct (nullable = true)
     |    |    |-- carName: string (nullable = true)
     |    |    |-- carNumbers: string (nullable = true)
     |    |    |-- driver: string (nullable = true)
     |-- teamName: array (nullable = true)
     |    |-- element: struct (containsNull = true)
     |    |    |-- team1: string (nullable = true)
     |    |    |-- team2: string (nullable = true)
    
    
    scala> df.show(false)
    +---------------------------+------------------------------+
    |info                       |teamName                      |
    +---------------------------+------------------------------+
    |[[RB7, 33, Max Verstappen]]|[[Redbull, rb], [Monster, mt]]|
    +---------------------------+------------------------------+
    
    

    您正在寻找的命令

    scala> val dfOut = df.nestedMapColumn(inputColumnName = "teamName", outputColumnName = "nextElementInArray", expression = a => concat(a.getField("team1"), a.getField("team2")) )
    dfOut: org.apache.spark.sql.DataFrame = [info: struct<drivers: struct<carName: string, carNumbers: string ... 1 more field>>, teamName: array<struct<team1:string,team2:string,nextElementInArray:string>>]
    

    输出

    scala> dfOut.printSchema
    root
     |-- info: struct (nullable = true)
     |    |-- drivers: struct (nullable = true)
     |    |    |-- carName: string (nullable = true)
     |    |    |-- carNumbers: string (nullable = true)
     |    |    |-- driver: string (nullable = true)
     |-- teamName: array (nullable = true)
     |    |-- element: struct (containsNull = false)
     |    |    |-- team1: string (nullable = true)
     |    |    |-- team2: string (nullable = true)
     |    |    |-- nextElementInArray: string (nullable = true)
    
    scala> dfOut.show(false)
    +---------------------------+----------------------------------------------------+
    |info                       |teamName                                            |
    +---------------------------+----------------------------------------------------+
    |[[RB7, 33, Max Verstappen]]|[[Redbull, rb, Redbullrb], [Monster, mt, Monstermt]]|
    +---------------------------+----------------------------------------------------+
    

    【讨论】:

      猜你喜欢
      • 2019-04-03
      • 2014-01-24
      • 1970-01-01
      • 2015-11-20
      • 1970-01-01
      • 2016-06-08
      • 2016-11-29
      • 2018-08-18
      • 2018-09-22
      相关资源
      最近更新 更多