【问题标题】:pyspark/dataframe - creating a nested structurepyspark/dataframe - 创建嵌套结构
【发布时间】:2023-03-19 06:02:01
【问题描述】:

我正在使用带有数据框的 pyspark,并希望创建一个嵌套结构,如下所示

之前:

Column 1 | Column 2 | Column 3 
--------------------------------
A    | B   | 1 
A    | B   | 2 
A    | C   | 1 

之后:

Column 1 | Column 4 
--------------------------------
A    | [B : [1,2]] 
A    | [C : [1]]

这可行吗?

【问题讨论】:

  • 是的,这是可行的。

标签: pyspark apache-spark-sql


【解决方案1】:

第一个可重现的数据框示例。

js = [{"col1": "A", "col2":"B", "col3":1},{"col1": "A", "col2":"B", "col3":2},{"col1": "A", "col2":"C", "col3":1}]
jsrdd = sc.parallelize(js)
sqlContext = SQLContext(sc)
jsdf = sqlContext.read.json(jsrdd)
jsdf.show()
+----+----+----+
|col1|col2|col3|
+----+----+----+
|   A|   B|   1|
|   A|   B|   2|
|   A|   C|   1|
+----+----+----+

现在,列表不存储为键值对。在 column2 上执行 groupby 后,您可以使用 dictionary 或简单的 collect_list()

jsdf.groupby(['col1', 'col2']).agg(F.collect_list('col3')).show()
+----+----+------------------+
|col1|col2|collect_list(col3)|
+----+----+------------------+
|   A|   C|               [1]|
|   A|   B|            [1, 2]|
+----+----+------------------+

【讨论】:

  • @undefined_variable 请编辑它以获得预期的输出,您可以在其中以键值对方式保存列表。谢谢。
【解决方案2】:

我认为您无法获得准确的输出,但您可以接近。问题是第 4 列的键名。在 Spark 中,结构需要预先知道一组固定的列。但是让我们稍后再说,首先是聚合:

import pyspark
from pyspark.sql import functions as F

sc = pyspark.SparkContext()
spark = pyspark.sql.SparkSession(sc)

data = [('A', 'B', 1), ('A', 'B', 2), ('A', 'C', 1)]
columns = ['Column1', 'Column2', 'Column3']

data = spark.createDataFrame(data, columns)

data.createOrReplaceTempView("data")
data.show()

# Result
+-------+-------+-------+
|Column1|Column2|Column3|
+-------+-------+-------+
|      A|      B|      1|
|      A|      B|      2|
|      A|      C|      1|
+-------+-------+-------+

nested = spark.sql("SELECT Column1, Column2, STRUCT(COLLECT_LIST(Column3) AS data) AS Column4 FROM data GROUP BY Column1, Column2")
nested.toJSON().collect()

# Result
['{"Column1":"A","Column2":"C","Column4":{"data":[1]}}',
 '{"Column1":"A","Column2":"B","Column4":{"data":[1,2]}}']

几乎是你想要的,对吧?问题是如果您事先不知道您的键名(即第 2 列中的值),Spark 无法确定您的数据结构。另外,我不完全确定如何将列的值用作结构的键,除非您使用 UDF(可能使用 PIVOT?):

datatype = 'struct<B:array<bigint>,C:array<bigint>>'  # Add any other potential keys here.
@F.udf(datatype)
def replace_struct_name(column2_value, column4_value):
    return {column2_value: column4_value['data']}

nested.withColumn('Column5', replace_struct_name(F.col("Column2"), F.col("Column4"))).toJSON().collect()

# Output
['{"Column1":"A","Column2":"C","Column4":{"C":[1]}}',
 '{"Column1":"A","Column2":"B","Column4":{"B":[1,2]}}']

这当然有一个缺点,键的数量必须是离散的,并且是事先知道的,否则其他键值将被默默地忽略。

【讨论】:

  • 谢谢,这很有帮助。在这种特殊情况下,我确实知道第 4 列的关键,所以它成功了。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2019-08-19
  • 1970-01-01
  • 1970-01-01
  • 2020-12-12
  • 1970-01-01
  • 1970-01-01
  • 2018-01-13
相关资源
最近更新 更多