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

PySpark中如何按指定item_name顺序对timestamp分组聚合?

问题

现有包含item_name、item_value、timestamp字段的PySpark DataFrame,当前执行以下代码按timestamp分组聚合item_value:

grouped_df = (
    df.groupBy("timestamp")
    .agg(sort_array(collect_list("item_values")).alias("item_values"))
)
result = [row["item_values"] for row in grouped_df.sort("timestamp").collect()]

但当前返回的item_value是按数值排序的,需要修改代码,让分组后的item_value严格按照asset_min、hpc_max、off_median的item_name固定顺序排列。

解决方案

实现思路

通过给目标item_name分配固定排序权重,在分组聚合时先按权重排序,再提取对应item_value,以此保证指定顺序。

修改后的代码

from pyspark.sql import functions as F

# 给需要固定顺序的item_name分配权重
df_with_weight = df.withColumn(
    "sort_weight",
    F.when(F.col("item_name") == "asset_min", 0)
     .when(F.col("item_name") == "hpc_max", 1)
     .when(F.col("item_name") == "off_median", 2)
     .otherwise(3)  # 其他未指定的item_name排在最后
)

# 分组后按权重排序并提取item_value
grouped_df = df_with_weight.groupBy("timestamp").agg(
    F.transform(
        F.sort_array(F.collect_list(F.struct("sort_weight", "item_value")), asc=True),
        lambda x: x.item_value
    ).alias("item_values")
)

# 获取最终结果
result = [row["item_values"] for row in grouped_df.sort("timestamp").collect()]

代码说明

  1. 添加权重列:用when函数为每个目标item_name分配对应的权重值,确保排序优先级符合需求;
  2. 分组聚合排序:分组时收集包含权重和item_value的结构体,通过sort_array按权重升序排列,再用transform提取出排序后的item_value;
  3. 结果输出:和原逻辑一致,按timestamp排序后收集结果。

这样处理后,每个分组内的item_value会严格按照asset_min → hpc_max → off_median的顺序排列,未指定的item_name会统一排在最后。

内容的提问来源于stack exchange,提问作者MMV

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.27 01:32:28