Spark SQL如何仅用一次GROUP BY实现日/周/月多维度分组统计
一次GROUP BY实现日/周/月多维度统计
当然可以通过一次GROUP BY完成多维度统计,核心思路是先将单条记录扩展为对应日、周、月三个维度的多条记录,再进行一次分组聚合,避免重复代码。以下是Spark SQL和PySpark两种实现方式:
Spark SQL 实现
WITH date_dimensions AS ( SELECT num, -- 将单条记录拆分为日、周、月三个维度的结构并展开 EXPLODE(ARRAY( STRUCT('day' AS dim_type, DATE(date) AS dim_value), STRUCT('week' AS dim_type, DATE_TRUNC('week', date) AS dim_value), STRUCT('month' AS dim_type, DATE_TRUNC('month', date) AS dim_value) )) AS dim FROM example ) SELECT dim.dim_type, dim.dim_value, SUM(num) AS total_num FROM date_dimensions GROUP BY dim.dim_type, dim.dim_value ORDER BY dim.dim_type, dim.dim_value;
- 逻辑说明:通过
EXPLODE函数将每条原始记录拆分为3条对应不同统计维度的记录,之后仅需一次GROUP BY即可完成所有维度的num求和,完全替代多次GROUP BY加UNION ALL的写法。
PySpark DataFrame 实现
from pyspark.sql import SparkSession from pyspark.sql.functions import sum, explode, array, struct, date_trunc, lit, col # 初始化SparkSession(若已存在可跳过) spark = SparkSession.builder.appName("MultiDimAgg").getOrCreate() # 加载example表为DataFrame df = spark.table("example") # 生成多维度记录并聚合统计 result_df = ( df.select( col("num"), explode( array( # 日维度:保留原始日期的日期部分 struct(lit("day").alias("dim_type"), col("date").cast("date").alias("dim_value")), # 周维度:截断到周起始日(Spark默认周一为周起始,可通过配置修改) struct(lit("week").alias("dim_type"), date_trunc("week", col("date")).cast("date").alias("dim_value")), # 月维度:截断到当月第一天 struct(lit("month").alias("dim_type"), date_trunc("month", col("date")).cast("date").alias("dim_value")) ) ).alias("dim") ) .select("dim.dim_type", "dim.dim_value", "num") .groupBy("dim_type", "dim_value") .agg(sum("num").alias("total_num")) .orderBy("dim_type", "dim_value") ) # 查看结果 result_df.show()
- 逻辑说明:利用PySpark的
array和struct构造多维度结构,通过explode展开后进行分组聚合,代码结构清晰,无重复逻辑。
内容的提问来源于stack exchange,提问作者Guoran Yun
相关产品推荐
相关产品推荐

