You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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)

代码说明

  1. Union操作:将两个DataFrame合并,保留所有行数据;
  2. Explode数组:把data数组展开成单行,便于对数组内的结构体进行聚合;
  3. 分组求和:先按主id、data_id、data_name分组,对seconds求和;同时通过窗口函数计算每个主id对应的total_seconds总和;
  4. 重组数组:将聚合后的结构体重新组合成data数组,最终得到目标格式。

内容的提问来源于stack exchange,提问作者Scott

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.23 11:54:55