PySpark同Schema DataFrame按ID连接并聚合结构体数组需求
PySpark DataFrame 合并与聚合实现
问题描述
现有两个Schema完全一致的PySpark DataFrame,Schema定义如下:
root |-- id: string (nullable = true) |-- data: array (nullable = true) | |-- element: struct (containsNull = true) | | |-- id: string (nullable = true) | | |-- name: string (nullable = true) | | |-- seconds: decimal(38,18) (nullable = true) |-- total_seconds: decimal(38,3) (nullable = true)
第一个DataFrame数据:
[{ "id": 1, "data": [{ "id": "123", "name": "name1", "seconds": 50 }, { "id": "234", "name": "name2", "seconds": 25 }], "total_seconds": 100 }, { "id": 2, "data": [{ "id": "123", "name": "name1", "seconds": 100 }, { "id": "234", "name": "name2", "seconds": 200 }], "total_seconds": 400 }]
第二个DataFrame数据:
[{ "id": 1, "data": [{ "id": "123", "name": "name1", "seconds": 100 }, { "id": "345", "name": "name3", "seconds": 25 }], "total_seconds": 400 }, { "id": 3, "data": [{ "id": "123", "name": "name1", "seconds": 50 }, { "id": "234", "name": "name2", "seconds": 100 }], "total_seconds": 200 }]
期望输出:
[{ "id": 1, "data": [{ "id": "123", "name": "name1", "seconds": 150 }, { "id": "234", "name": "name2", "seconds": 25 }, { "id": "345", "name": "name3", "seconds": 25 }], "total_seconds": 500 }, { "id": 2, "data": [{ "id": "123", "name": "name1", "seconds": 100 }, { "id": "234", "name": "name2", "seconds": 200 }], "total_seconds": 400 }, { "id": 3, "data": [{ "id": "123", "name": "name1", "seconds": 50 }, { "id": "234", "name": "name2", "seconds": 100 }], "total_seconds": 200 }]
需要完成以下操作:
- 按id列进行合并
- 聚合total_seconds字段(求和)
- 聚合/合并data列:对同一id下的结构体数组,按data.id和data.name合并,seconds字段求和,最终重新组合为结构体数组
解决方案
通过联合两个DataFrame + 分组聚合的方式实现需求,具体代码如下:
1. 初始化DataFrame
首先创建两个示例DataFrame:
from pyspark.sql import SparkSession from pyspark.sql.types import ( StructType, StructField, StringType, ArrayType, DecimalType ) from pyspark.sql import functions as F from pyspark.sql.window import Window # 初始化SparkSession spark = SparkSession.builder.appName("DataMerge").getOrCreate() # 定义Schema schema = StructType([ StructField("id", StringType(), nullable=True), StructField("data", ArrayType( StructType([ StructField("id", StringType(), nullable=True), StructField("name", StringType(), nullable=True), StructField("seconds", DecimalType(38, 18), nullable=True) ]), containsNull=True ), nullable=True), StructField("total_seconds", DecimalType(38, 3), nullable=True) ]) # 第一个DataFrame数据 df1_data = [ ("1", [("123", "name1", 50), ("234", "name2", 25)], 100), ("2", [("123", "name1", 100), ("234", "name2", 200)], 400) ] df1 = spark.createDataFrame(df1_data, schema) # 第二个DataFrame数据 df2_data = [ ("1", [("123", "name1", 100), ("345", "name3", 25)], 400), ("3", [("123", "name1", 50), ("234", "name2", 100)], 200) ] df2 = spark.createDataFrame(df2_data, schema)
2. 合并并聚合数据
# 联合两个DataFrame union_df = df1.union(df2) # 展开data数组,方便后续聚合 exploded_df = union_df.withColumn("data_element", F.explode("data")) # 提取结构体中的字段 flattened_df = exploded_df.select( "id", F.col("data_element.id").alias("data_id"), F.col("data_element.name").alias("data_name"), F.col("data_element.seconds").alias("data_seconds"), "total_seconds" ) # 分组聚合:按主id、data_id、data_name求和seconds,同时按主id求和total_seconds aggregated_df = flattened_df.groupBy("id", "data_id", "data_name").agg( F.sum("data_seconds").alias("seconds"), F.first(F.sum("total_seconds").over(Window.partitionBy("id"))).alias("total_seconds") ) # 将聚合后的data字段重新组合成数组 result_df = aggregated_df.groupBy("id", "total_seconds").agg( F.collect_list(F.struct("data_id", "data_name", "seconds")).alias("data") ).orderBy("id") # 查看结果 result_df.show(truncate=False)
代码说明
- Union操作:将两个DataFrame合并,保留所有行数据;
- Explode数组:把
data数组展开成单行,便于对数组内的结构体进行聚合; - 分组求和:先按主
id、data_id、data_name分组,对seconds求和;同时通过窗口函数计算每个主id对应的total_seconds总和; - 重组数组:将聚合后的结构体重新组合成
data数组,最终得到目标格式。
内容的提问来源于stack exchange,提问作者Scott
相关产品推荐
相关产品推荐

