【问题标题】:Pyspark: How to flatten nested arrays by merging values in sparkPyspark:如何通过合并 Spark 中的值来展平嵌套数组
【发布时间】:2021-10-20 13:52:38
【问题描述】:

我有 10000 个具有不同 ID 的 json,每个都有 10000 个名称。如何通过在pyspark中通过int或str合并值来展平嵌套数组?

编辑:我添加了列name_10000_xvz 来解释更好的数据结构。我也更新了 Notes、Input df、所需的输出 df 和输入 json 文件。

注意事项:

  • 输入数据帧有超过 10000 个列 name_1_a、name_1000_xx 所以列(数组)名称不能硬编码,因为它需要写入 10000 个名称
  • iddateval 在所有列和所有 json 中始终具有相同的命名约定
  • 数组大小可以变化,但dateval 始终存在,因此可以硬编码
  • date 在每个数组中可以不同,例如 name_1_a 以 2001 开头,但 name_10000_xvz for id == 1 从 2000 开始,finnish 从 2004 开始,但是对于 id == 2,从 1990 开始并以 2004 结束

输入df:

root
 |-- id: long (nullable = true)
 |-- name_10000_xvz: array (nullable = true)
 |    |-- element: struct (containsNull = true)
 |    |    |-- date: long (nullable = true)
 |    |    |-- val: long (nullable = true)
 |-- name_1_a: array (nullable = true)
 |    |-- element: struct (containsNull = true)
 |    |    |-- date: long (nullable = true)
 |    |    |-- val: long (nullable = true)
 |-- name_1_b: array (nullable = true)
 |    |-- element: struct (containsNull = true)
 |    |    |-- date: long (nullable = true)
 |    |    |-- val: long (nullable = true)
 |-- name_2_a: array (nullable = true)
 |    |-- element: struct (containsNull = true)
 |    |    |-- date: long (nullable = true)
 |    |    |-- val: long (nullable = true)

+---+------------------------------------------------------------------------+---------------------------------+---------------------------------+------------------------------------+
|id |name_10000_xvz                                                          |name_1_a                         |name_1_b                         |name_2_a                            |
+---+------------------------------------------------------------------------+---------------------------------+---------------------------------+------------------------------------+
|2  |[{1990, 39}, {2000, 30}, {2001, 31}, {2002, 32}, {2003, 33}, {2004, 34}]|[{2001, 1}, {2002, 2}, {2003, 3}]|[{2001, 4}, {2002, 5}, {2003, 6}]|[{2001, 21}, {2002, 22}, {2003, 23}]|
|1  |[{2000, 30}, {2001, 31}, {2002, 32}, {2003, 33}]                        |[{2001, 1}, {2002, 2}, {2003, 3}]|[{2001, 4}, {2002, 5}, {2003, 6}]|[{2001, 21}, {2002, 22}, {2003, 23}]|
+---+------------------------------------------------------------------------+---------------------------------+---------------------------------+------------------------------------+

所需的输出df:

+---+---------+----------+-----------+---------+----------------+
|id |   date  | name_1_a | name_1_b  |name_2_a | name_10000_xvz |
+---+---------+----------+-----------+---------+----------------+
|1  |   2000  |     0    |    0      |   0     |        30      |
|1  |   2001  |     1    |    4      |   21    |        31      |
|1  |   2002  |     2    |    5      |   22    |        32      |
|1  |   2003  |     3    |    6      |   23    |        33      |
|2  |   1990  |     0    |    0      |   0     |        39      |
|2  |   2000  |     0    |    0      |   0     |        30      |
|2  |   2001  |     1    |    4      |   21    |        31      |
|2  |   2002  |     2    |    5      |   22    |        32      |
|2  |   2003  |     3    |    6      |   23    |        33      |
|2  |   2004  |     0    |    0      |   0     |        34      |
+---+---------+----------+-----------+---------+----------------+

重现输入 df:

df = spark.read.json(sc.parallelize([
  """{"id":1,"name_1_a":[{"date":2001,"val":1},{"date":2002,"val":2},{"date":2003,"val":3}],"name_1_b":[{"date":2001,"val":4},{"date":2002,"val":5},{"date":2003,"val":6}],"name_2_a":[{"date":2001,"val":21},{"date":2002,"val":22},{"date":2003,"val":23}],"name_10000_xvz":[{"date":2000,"val":30},{"date":2001,"val":31},{"date":2002,"val":32},{"date":2003,"val":33}]}""",
  """{"id":2,"name_1_a":[{"date":2001,"val":1},{"date":2002,"val":2},{"date":2003,"val":3}],"name_1_b":[{"date":2001,"val":4},{"date":2002,"val":5},{"date":2003,"val":6}],"name_2_a":[{"date":2001,"val":21},{"date":2002,"val":22},{"date":2003,"val":23}],"name_10000_xvz":[{"date":1990,"val":39},{"date":2000,"val":30},{"date":2001,"val":31},{"date":2002,"val":32},{"date":2003,"val":33},{"date":2004,"val":34}]}}"""
]))

有用的链接:

【问题讨论】:

    标签: python sql apache-spark pyspark apache-spark-sql


    【解决方案1】:

    更新

    正如@werner 所提到的,有必要转换所有结构以将列名附加到其中。

    import pyspark.sql.functions as f
    
    names = [column for column in df.columns if column.startswith('name_')]
    
    expressions = []
    for name in names:
      expressions.append(f.expr('TRANSFORM({name}, el -> STRUCT("{name}" AS name, el.date, el.val))'.format(name=name)))
    
    flatten_df = (df
                  .withColumn('flatten', f.flatten(f.array(*expressions)))
                  .selectExpr('id', 'inline(flatten)'))
    
    output_df = (flatten_df
                 .groupBy('id', 'date')
                 .pivot('name', names)
                 .agg(f.first('val')))
    
    output_df.sort('id', 'date').show(truncate=False)
    +---+----+--------------+--------+--------+--------+
    |id |date|name_10000_xvz|name_1_a|name_1_b|name_2_a|
    +---+----+--------------+--------+--------+--------+
    |1  |2000|30            |null    |null    |null    |
    |1  |2001|31            |1       |4       |21      |
    |1  |2002|32            |2       |5       |22      |
    |1  |2003|33            |3       |6       |23      |
    |2  |1990|39            |null    |null    |null    |
    |2  |2000|30            |null    |null    |null    |
    |2  |2001|31            |1       |4       |21      |
    |2  |2002|32            |2       |5       |22      |
    |2  |2003|33            |3       |6       |23      |
    |2  |2004|34            |null    |null    |null    |
    +---+----+--------------+--------+--------+--------+
    

    假设:

    • date 值始终是所有列的相同值
    • name_1_a, name_1_b, name_2_a 它们的大小相等
    import pyspark.sql.functions as f
    
    output_df = (df
                 .withColumn('flatten', f.expr('TRANSFORM(SEQUENCE(0, size(name_1_a) - 1), i -> ' \
                                               'STRUCT(name_1_a[i].date AS date, ' \
                                               '       name_1_a[i].val AS name_1_a, ' \
                                               '       name_1_b[i].val AS name_1_b, ' \
                                               '       name_2_a[i].val AS name_2_a))'))
                 .selectExpr('id', 'inline(flatten)'))
    
    output_df.sort('id', 'date').show(truncate=False)
    +---+----+--------+--------+--------+
    |id |date|name_1_a|name_1_b|name_2_a|
    +---+----+--------+--------+--------+
    |1  |2001|1       |4       |21      |
    |1  |2002|2       |5       |22      |
    |1  |2003|3       |6       |23      |
    |2  |2001|1       |4       |21      |
    |2  |2002|2       |5       |22      |
    |2  |2003|3       |6       |23      |
    +---+----+--------+--------+--------+
    

    【讨论】:

    • 感谢您的帮助,但名称不能硬编码,因为您的解决方案需要写入 10000 个名称,数组大小可能会有所不同,我已经更新了问题
    • 之前不知道inline
    • 非常感谢,你是明星,我又用类似的数据提出了一个问题,如果你能提供帮助,那就太好了。 stackoverflow.com/questions/68896130/…
    【解决方案2】:

    如何使用命名约定?

    你能用 spark-sql 试试下面的东西吗?

    df.createOrReplaceTempView("df")
    spark.sql("""
    select id, 
    name_1_a.date[0] as date, name_1_a.val[0] as name_1_a, name_1_b.val[0] as name_1_b, name_2_a.val[0] as name_2_a
    from df
    """).show(false)
    
    +---+----+--------+--------+--------+
    |id |date|name_1_a|name_1_b|name_2_a|
    +---+----+--------+--------+--------+
    |1  |2001|1       |4       |21      |
    |2  |2001|1       |4       |21      |
    +---+----+--------+--------+--------+
    

    这是我的假设。

    1. 第一个字段是 id,其余都是名称..1 到 n,例如 name_1_a、name_1_b、name_2_a 等
    2. 所有“n”个名称的日期都相同,因此我可以使用第一个字段来推导它。

    构建数据框。

    JSON 字符串

    val jsonstr1 = """{  "id": 1,  "name_1_a": [    {      "date": 2001,      "val": 1    },    {      "date": 2002,      "val": 2    },    {      "date": 2003,      "val": 3    }  ],  "name_1_b": [    {      "date": 2001,      "val": 4    },    {      "date": 2002,      "val": 5    },    {      "date": 2003,      "val": 6    }  ],  "name_2_a": [    {      "date": 2001,      "val": 21    },    {      "date": 2002,      "val": 22    },    {      "date": 2003,      "val": 23    }  ]}"""
    
    val jsonstr2 = """{  "id": 2,  "name_1_a": [    {      "date": 2001,      "val": 1    },    {      "date": 2002,      "val": 2    },    {      "date": 2003,      "val": 3    }  ],  "name_1_b": [    {      "date": 2001,      "val": 4    },    {      "date": 2002,      "val": 5    },    {      "date": 2003,      "val": 6    }  ],  "name_2_a": [    {      "date": 2001,      "val": 21    },    {      "date": 2002,      "val": 22    },    {      "date": 2003,      "val": 23    }  ]}"""
    

    数据帧

    val df1 = spark.read.json(Seq(jsonstr1).toDS)
    val df2 = spark.read.json(Seq(jsonstr2).toDS)
    val df = df1.union(df2)
    

    现在在 df 上创建一个视图。我只是将其命名为“df”

    df.createOrReplaceTempView("df")
    

    显示数据:

    df.show(false)
    df.printSchema
    

    使用元数据并构造 sql 字符串。

    df.columns
    

    Array[String] = Array(id, name_1_a, name_1_b, name_2_a)

    val names = df.columns.drop(1)  // drop id
    val sql1 = for { i <- 0 to 2
                      t1=names.map( x => x + s".val[${i}] as ${x}").mkString(",")
                      t2 = names(0) + ".date[0] as date ," + t1
                  _=println(t)
        } yield s""" select id, ${t2} from df """
    val sql2 = sql1.mkString(" union All ")
    

    现在 sql2 包含以下字符串,这是一个有效的 sql

    " select id, name_1_a.date[0] as date ,name_1_a.val[0] as name_1_a,name_1_b.val[0] as name_1_b,name_2_a.val[0] as name_2_a from df  union All  select id, name_1_a.date[0] as date ,name_1_a.val[1] as name_1_a,name_1_b.val[1] as name_1_b,name_2_a.val[1] as name_2_a from df  union All  select id, name_1_a.date[0] as date ,name_1_a.val[2] as name_1_a,name_1_b.val[2] as name_1_b,name_2_a.val[2] as name_2_a from df "
    

    将其传递给 spark.sql(sql2) 并获得所需的结果

    spark.sql(sql2).orderBy("id").show(false)
    
    +---+----+--------+--------+--------+
    |id |date|name_1_a|name_1_b|name_2_a|
    +---+----+--------+--------+--------+
    |1  |2001|2       |5       |22      |
    |1  |2001|1       |4       |21      |
    |1  |2001|3       |6       |23      |
    |2  |2001|1       |4       |21      |
    |2  |2001|3       |6       |23      |
    |2  |2001|2       |5       |22      |
    +---+----+--------+--------+--------+
    

    【讨论】:

    • 对列名进行短语化(因为它们不能被硬编码)df.columns 并用于 sql 查询是个好主意。如果日期在所有“n”个名称中变化,sql 查询将如何?我已添加列 name_10000_xvz 以更详细地解释数据结构。
    猜你喜欢
    • 1970-01-01
    • 2023-03-03
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2016-02-26
    • 2023-03-03
    • 1970-01-01
    相关资源
    最近更新 更多