将Pandas的assign_boxes函数转换为PySpark等效函数求助
PySpark 实现等效的
assign_boxes 逻辑 先明确原Pandas代码的核心逻辑:
- 按
store分组,计算每组boxes的总和total - 计算
d:取total // 100与组内记录数-1的最小值 - 生成对应长度的列表:前
d个元素为100,第d+1个元素为total - 100*d,剩余元素补0,最终将这个列表按组内顺序映射回原数据
实现步骤
- 给每组内的记录添加行号,保证后续能按原顺序分配值(这里用
worker字段作为组内排序依据,对应原Pandas的组内默认顺序) - 计算每组的
total、组内记录数count,进而推导d和剩余值remainder = total - 100*d - 用窗口函数结合行号,判断每条记录对应的
optimal_boxes值
完整代码
from pyspark.sql import functions as F from pyspark.sql.window import Window # 假设你已创建df_stack_exchange,以下是核心实现逻辑 # 步骤1:给每组内记录按worker排序并添加行号 window_group = Window.partitionBy("store").orderBy("worker") df_with_row_num = df_stack_exchange.withColumn("row_num", F.row_number().over(window_group)) # 步骤2:计算每组的聚合参数 group_agg = df_with_row_num.groupBy("store").agg( F.sum("boxes").alias("total"), F.count("*").alias("count") ).withColumn("d", F.least(F.floor(F.col("total") / 100), F.col("count") - 1)) \ .withColumn("remainder", F.col("total") - F.col("d") * 100) # 步骤3:关联聚合结果,生成optimal_boxes字段 result_df = df_with_row_num.join(group_agg, on="store", how="left") \ .withColumn("optimal_boxes", F.when(F.col("row_num") <= F.col("d"), 100) .when(F.col("row_num") == F.col("d") + 1, F.col("remainder")) .otherwise(0) ) \ .drop("row_num", "total", "count", "d", "remainder") # 查看结果 result_df.show()
结果验证
运行后输出与原Pandas代码结果完全一致:
| store | worker | boxes | optimal_boxes |
|---|---|---|---|
| A | 1 | 105 | 105 |
| B | 1 | 90 | 100 |
| B | 2 | 100 | 90 |
| C | 1 | 80 | 100 |
| C | 2 | 10 | 100 |
| C | 3 | 200 | 90 |
| D | 1 | 70 | 100 |
| D | 2 | 210 | 100 |
| D | 3 | 50 | 100 |
| D | 4 | 0 | 30 |
内容的提问来源于stack exchange,提问作者qwerty12
相关产品推荐
相关产品推荐

