Spark中嵌套groupBy与agg实现扁平表转指定struct嵌套结构
实现方案
核心是两层分组聚合配合嵌套结构构造,可通过Spark SQL/PySpark实现,步骤如下:
实现逻辑
- 第一层聚合:按
transaction_id、transaction_date、partition_key三个字段分组,计算每个分区的金额总和,同时收集该分区下所有的记录明细,生成嵌套的records结构 - 第二层聚合:按
transaction_id、transaction_date两个字段分组,将同一交易下不同分区的结果收集为transaction_partition数组 - 字段类型对齐:按要求将对应字段强转为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
相关产品推荐
相关产品推荐

