如何使用PySpark转置SparkSQL聚合查询输出的DataFrame结果
实现方案
假设你执行SQL后得到的单行聚合DataFrame名为agg_df,可按如下方式实现转置:
方法1:固定字段写法(适合字段少的场景)
首先导入依赖:
from pyspark.sql.functions import expr
然后执行转置逻辑:
transposed_df = agg_df.select(expr(""" stack( 2, 'age', `max(age)`, `min(age)`, `avg(age)`, 'sal', `max(sal)`, `min(sal)`, `avg(sal)` ) AS (columns, max, min, avg) """))
说明:
- stack函数第一个参数为要拆分的行数,这里对应age、sal共2行
- 每行的参数依次为:字段名、该字段对应的max聚合值、min聚合值、avg聚合值
- 原DataFrame的列名包含特殊字符,需要用反引号
`包裹避免语法报错
方法2:动态生成写法(适合字段多的场景)
如果后续需要新增统计字段,可通过代码动态生成stack表达式,避免手动修改:
from pyspark.sql.functions import expr # 配置要统计的字段列表 stat_fields = ["age", "sal"] # 自动拼接stack参数 stack_params = [] for field in stat_fields: stack_params.append(f"'{field}', `max({field})`, `min({field})`, `avg({field})`") # 生成完整的stack表达式 stack_expr = f"stack({len(stat_fields)}, {','.join(stack_params)}) AS (columns, max, min, avg)" # 执行转置 transposed_df = agg_df.select(expr(stack_expr))
验证结果
执行transposed_df.show()后输出符合要求的结果:
+-------+-----+----+----+ |columns| max| min| avg| +-------+-----+----+----+ | age| 46| 23| 31| | sal|10000|2000|5000| +-------+-----+----+----+
内容的提问来源于stack exchange,提问作者Satyabrata
相关产品推荐
相关产品推荐

