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

Spark中嵌套groupBy与agg实现扁平表转指定struct嵌套结构

实现方案

核心是两层分组聚合配合嵌套结构构造,可通过Spark SQL/PySpark实现,步骤如下:

实现逻辑

  1. 第一层聚合:按transaction_id、transaction_date、partition_key三个字段分组,计算每个分区的金额总和,同时收集该分区下所有的记录明细,生成嵌套的records结构
  2. 第二层聚合:按transaction_id、transaction_date两个字段分组,将同一交易下不同分区的结果收集为transaction_partition数组
  3. 字段类型对齐:按要求将对应字段强转为String类型,最后生成指定格式的JSON结构即可

代码示例

Spark SQL 写法

假设源表名为transaction_source,代码如下:

WITH partition_level_agg AS (
    SELECT
        transaction_id,
        transaction_date,
        partition_key,
        SUM(amount) AS record_amount_sum,
        COLLECT_LIST(
            STRUCT(
                CAST(record_id AS STRING) AS record_id,
                amount AS record_amount,
                CAST(record_in_date AS STRING) AS record_in_date
            )
        ) AS records
    FROM transaction_source
    GROUP BY transaction_id, transaction_date, partition_key
)
SELECT
    CAST(transaction_id AS STRING) AS transaction_id,
    TO_JSON(
        STRUCT(
            transaction_id,
            transaction_date,
            COLLECT_LIST(
                STRUCT(
                    partition_key,
                    record_amount_sum,
                    records
                )
            ) AS transaction_partition
        )
    ) AS transaction
FROM partition_level_agg
GROUP BY transaction_id, transaction_date

PySpark DataFrame 写法

from pyspark.sql import functions as F
from pyspark.sql.types import StringType

# 源数据表为source_df
partition_agg = source_df.groupBy("transaction_id", "transaction_date", "partition_key") \
    .agg(
        F.sum("amount").alias("record_amount_sum"),
        F.collect_list(
            F.struct(
                F.col("record_id").cast(StringType()).alias("record_id"),
                F.col("amount").alias("record_amount"),
                F.col("record_in_date").cast(StringType()).alias("record_in_date")
            )
        ).alias("records")
    )

result_df = partition_agg.groupBy("transaction_id", "transaction_date") \
    .agg(
        F.collect_list(
            F.struct(
                "partition_key",
                "record_amount_sum",
                "records"
            )
        ).alias("transaction_partition")
    ) \
    .select(
        F.col("transaction_id").cast(StringType()).alias("transaction_id"),
        F.to_json(F.struct("transaction_id", "transaction_date", "transaction_partition")).alias("transaction")
    )

# 输出验证结果
result_df.show(truncate=False)

注意事项

  • 所有要求为String类型的字段都主动做了类型转换,避免源表字段类型不匹配Schema要求
  • 如果使用Hive SQL实现,语法逻辑和上述Spark SQL基本一致,仅需注意Hive版本对to_json函数的兼容性,低版本可手动拼接JSON字符串替代
  • 如果需要输出的是嵌套结构而非JSON字符串,去掉to_json函数的调用即可直接得到对应Schema的结构化数据

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.03 08:39:02