如何在PySpark透视表中处理多sum值并完成后续聚合计算?
解决PySpark透视表后计算与聚合的问题
步骤1:明确透视后的列结构
执行groupby("stock", "hour").pivot("day").sum("price_x_units", "units")后,生成的pivotdf列名会以[日期]_[字段名]的格式呈现,比如日期为20240101的话,对应列就是20240101_price_x_units和20240101_units。
步骤2:批量计算每个日期的均价(price_x_units / units)
先提取所有透视后的日期值,再对每个日期对应的两列做除法运算,同时处理除数为0的异常情况:
from pyspark.sql.functions import col, when # 提取所有唯一的日期值 days = list(set([col.split("_")[0] for col in pivotdf.columns if "_price_x_units" in col])) # 生成每个日期的均价计算表达式 calc_exprs = [] for day in days: # 避免除以0,这里将units为0的情况设为null,可根据业务调整为0或其他值 avg_price_col = when(col(f"{day}_units") != 0, col(f"{day}_price_x_units") / col(f"{day}_units")).alias(f"{day}_avg_price") calc_exprs.append(avg_price_col) # 保留stock、hour字段,加上计算出的均价列 avg_df = pivotdf.select("stock", "hour", *calc_exprs)
步骤3:去除stock维度并聚合汇总
按hour分组,对每个日期的均价列执行聚合操作(这里以均值为例,可根据需求替换为sum等函数):
from pyspark.sql.functions import avg # 生成聚合表达式 agg_exprs = [avg(col(f"{day}_avg_price")).alias(f"{day}_avg_price") for day in days] # 执行聚合得到最终汇总表 final_df = avg_df.groupby("hour").agg(*agg_exprs)
完整示例代码
假设原始数据样例:
+-----+--------+----+---------------+------+ |stock| day|hour|price_x_units|units| +-----+--------+----+---------------+------+ | AAPL|20240101| 9| 1000| 10| | AAPL|20240101| 9| 1500| 15| | MSFT|20240101| 9| 2000| 20| | AAPL|20240102| 10| 3000| 30| +-----+--------+----+---------------+------+
完整处理代码:
from pyspark.sql.functions import col, when, avg # 生成透视表 pivotdf = df.groupby("stock", "hour").pivot("day").sum("price_x_units", "units") # 提取所有日期值 days = list(set([col.split("_")[0] for col in pivotdf.columns if "_price_x_units" in col])) # 计算每个日期的均价 calc_exprs = [] for day in days: avg_price_col = when(col(f"{day}_units") != 0, col(f"{day}_price_x_units") / col(f"{day}_units")).alias(f"{day}_avg_price") calc_exprs.append(avg_price_col) avg_df = pivotdf.select("stock", "hour", *calc_exprs) # 按hour聚合去除stock维度 agg_exprs = [avg(col(f"{day}_avg_price")).alias(f"{day}_avg_price") for day in days] final_df = avg_df.groupby("hour").agg(*agg_exprs) final_df.show()
注意事项
- 务必处理除数为0的场景:如果
units为0,直接计算会返回null或报错,可通过when语句设置默认值。 - 聚合函数按需选择:如果需要汇总所有stock的均价总和,把
avg替换为sum即可,完全匹配业务需求。
内容的提问来源于stack exchange,提问作者zenith7
相关产品推荐
相关产品推荐

