【问题标题】:mininum value of struct type column in a dataframe数据框中结构类型列的最小值
【发布时间】:2019-01-17 12:27:33
【问题描述】:

我的问题是我有下面的 json 文件,其中包含 column3 的结构类型数据。我可以提取行但无法找到 column3 的最小值。其中 column3 包含带有值的动态嵌套列(动态名称)。

输入数据是:

"result": { "data" : 
[ {"col1": "value1",  "col2": "value2",  "col3" : { "dyno" : 3, "aeio": 5 }, "col4": "value4"},
   {"col1": "value11", "col2": "value22", "col3" : { "abc" : 6, "def": 9 , "aero": 2}, "col4": "value44"},
   {"col1": "value12", "col2": "value23", "col3" : { "ddc" : 6}, "col4": "value43"}]  

outputDate 预期为:

col1    col2    col3    col4    col5(min value of col3)

value1  value2  [3,5]   value4  3

value11 value22 [6,9,2] value44 2

value12 value23 [6] value43 6

我可以读取文件并将数据分解为记录,但无法找到 col3 的最小值。

val bestseller_df1 = bestseller_json.withColumn("extractedresult", explode(col("result.data")))

能否请您帮我编写代码以查找 spark/scala 中 col3 的最小值。

我的 json 文件是:

{"success":true, "result": { "data": [ {"col1": "value1",  "col2": "value2",  "col3" : { "dyno" : 3, "aeio": 5 }, "col4": "value4"},{"col1": "value11", "col2": "value22", "col3" : { "abc" : 6, "def": 9 , "aero": 2}, "col4": "value44"},{"col1": "value12", "col2": "value23", "col3" : { "ddc" : 6}, "col4": "value43"}],"total":3}}

【问题讨论】:

    标签: scala apache-spark


    【解决方案1】:

    你会怎么做

    scala> val df = spark.read.json("/tmp/stack/pathi.json")
    df: org.apache.spark.sql.DataFrame = [result: struct<data: array<struct<col1:string,col2:string,col3:struct<abc:bigint,aeio:bigint,aero:bigint,ddc:bigint,def:bigint,dyno:bigint>,col4:string>>, total: bigint>, success: boolean]
    
    scala> df.printSchema
    root
     |-- result: struct (nullable = true)
     |    |-- data: array (nullable = true)
     |    |    |-- element: struct (containsNull = true)
     |    |    |    |-- col1: string (nullable = true)
     |    |    |    |-- col2: string (nullable = true)
     |    |    |    |-- col3: struct (nullable = true)
     |    |    |    |    |-- abc: long (nullable = true)
     |    |    |    |    |-- aeio: long (nullable = true)
     |    |    |    |    |-- aero: long (nullable = true)
     |    |    |    |    |-- ddc: long (nullable = true)
     |    |    |    |    |-- def: long (nullable = true)
     |    |    |    |    |-- dyno: long (nullable = true)
     |    |    |    |-- col4: string (nullable = true)
     |    |-- total: long (nullable = true)
     |-- success: boolean (nullable = true)
    
    
    scala> df.show(false)
    +-------------------------------------------------------------------------------------------------------------------------------+-------+
    |result                                                                                                                         |success|
    +-------------------------------------------------------------------------------------------------------------------------------+-------+
    |[[[value1, value2, [, 5,,,, 3], value4], [value11, value22, [6,, 2,, 9,], value44], [value12, value23, [,,, 6,,], value43]], 3]|true   |
    +-------------------------------------------------------------------------------------------------------------------------------+-------+
    
    scala> df.select(explode($"result.data")).show(false)
    +-----------------------------------------+
    |col                                      |
    +-----------------------------------------+
    |[value1, value2, [, 5,,,, 3], value4]    |
    |[value11, value22, [6,, 2,, 9,], value44]|
    |[value12, value23, [,,, 6,,], value43]   |
    +-----------------------------------------+
    

    通过查看架构,现在我们知道“col3”中可能的列的列表,因此我们可以通过如下硬编码计算所有这些值的最小值

    scala> df.select(explode($"result.data")).select(least($"col.col3.abc",$"col.col3.aeio",$"col.col3.aero",$"col.col3.ddc",$"col.col3.def",$"col.col3.dyno")).show(false)
    +--------------------------------------------------------------------------------------------+
    |least(col.col3.abc, col.col3.aeio, col.col3.aero, col.col3.ddc, col.col3.def, col.col3.dyno)|
    +--------------------------------------------------------------------------------------------+
    |3                                                                                           |
    |2                                                                                           |
    |6                                                                                           |
    +--------------------------------------------------------------------------------------------+
    

    动态处理:

    我假设直到 col.col3,结构保持不变,所以我们继续创建另一个数据框作为

    scala> val df2 = df.withColumn("res_data",explode($"result.data")).select(col("success"),col("res_data"),$"res_data.col3.*")
    df2: org.apache.spark.sql.DataFrame = [success: boolean, res_data: struct<col1: string, col2: string ... 2 more fields> ... 6 more fields]
    
    scala> df2.show(false)
    +-------+-----------------------------------------+----+----+----+----+----+----+
    |success|res_data                                 |abc |aeio|aero|ddc |def |dyno|
    +-------+-----------------------------------------+----+----+----+----+----+----+
    |true   |[value1, value2, [, 5,,,, 3], value4]    |null|5   |null|null|null|3   |
    |true   |[value11, value22, [6,, 2,, 9,], value44]|6   |null|2   |null|9   |null|
    |true   |[value12, value23, [,,, 6,,], value43]   |null|null|null|6   |null|null|
    +-------+-----------------------------------------+----+----+----+----+----+----+
    

    除了“success”和“res_data”之外,其余列都是来自“col3”的动态列

    scala> val p = df2.columns
    p: Array[String] = Array(success, res_data, abc, aeio, aero, ddc, def, dyno)
    

    过滤这两个并将其余的映射到火花列

    scala> val m = p.filter(_!="success").filter(_!="res_data").map(col(_))
    m: Array[org.apache.spark.sql.Column] = Array(abc, aeio, aero, ddc, def, dyno)
    

    现在将m:_* 作为参数传递给最小函数,您会得到如下结果

    scala> df2.withColumn("minv",least(m:_*)).show(false)
    +-------+-----------------------------------------+----+----+----+----+----+----+----+
    |success|res_data                                 |abc |aeio|aero|ddc |def |dyno|minv|
    +-------+-----------------------------------------+----+----+----+----+----+----+----+
    |true   |[value1, value2, [, 5,,,, 3], value4]    |null|5   |null|null|null|3   |3   |
    |true   |[value11, value22, [6,, 2,, 9,], value44]|6   |null|2   |null|9   |null|2   |
    |true   |[value12, value23, [,,, 6,,], value43]   |null|null|null|6   |null|null|6   |
    +-------+-----------------------------------------+----+----+----+----+----+----+----+
    
    
    scala>
    

    希望这会有所帮助。

    【讨论】:

      【解决方案2】:

      dbutils.fs.put("/tmp/test.json", """

      {“col1”:“value1”,“col2”:“value2”,“col3”:{“dyno”:3,“aeio”:5},“col4”:“value4”},

      {“col1”:“value11”,“col2”:“value22”,“col3”:{“abc”:6,“def”:9,“aero”:2},“col4”:“value44 "},

      {“col1”:“value12”,“col2”:“value23”,“col3”:{“ddc”:6},“col4”:“value43”}]} """, 真)

      val df_json = spark.read.json("/tmp/test.json")

      val tf = df_json.withColumn("col3", explode(array($"col3.*"))).toDF

      val tmp_group = tf.groupBy("col1").agg( min(tf.col("col3")).alias("col3"))

      val top_rows = tf.join(tmp_group, Seq("col3","col1"), "inner")

      top_rows.select("col1", "col2", "col3","col4").show()

      写了 282 个字节。

      +-------+-------+----+-------+

      | col1| col2|col3| col4|

      +-------+-------+----+-------+

      |值1|值2| 3|值4|

      |值11|值22| 2|value44|

      |值12|值23| 6|值43|

      +-------+-------+----+-------+

      【讨论】:

        猜你喜欢
        • 2023-03-16
        • 1970-01-01
        • 1970-01-01
        • 2019-01-02
        • 2011-01-04
        • 1970-01-01
        • 1970-01-01
        • 2018-01-08
        • 1970-01-01
        相关资源
        最近更新 更多