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

如何使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.24 08:45:04