如何在PySpark中实现符合规则的ID动态分配
PySpark多门店ID分配逻辑实现方案
输入DataFrame定义
先明确三个核心输入表的结构,以下是示例数据:
1. store_df(门店已有ID映射)
from pyspark.sql import SparkSession from pyspark.sql import functions as F from pyspark.sql.window import Window spark = SparkSession.builder.appName("IDAllocation").getOrCreate() store_df = spark.createDataFrame([ ("StoreA", "ID1", "Div1"), ("StoreA", "ID3", "Div1"), ("StoreB", "ID2", "Div2"), ("StoreB", "ID5", "Div2") ], ["Store", "ID", "Div"])
2. final_list(待选ID池)
final_list = spark.createDataFrame([ ("Div1", "ID1", 1, "A"), ("Div1", "ID2", 2, "A"), ("Div1", "ID3", 3, "B"), ("Div1", "ID4", 4, "B"), ("Div1", "ID6", 5, "C"), ("Div2", "ID2", 1, "A"), ("Div2", "ID4", 2, "B"), ("Div2", "ID5", 3, "C"), ("Div2", "ID7", 4, "D"), ("Div2", "ID8", 5, "D") ], ["Div", "ID", "Rank", "Category"])
3. max_df(门店分配限额)
max_df = spark.createDataFrame([ ("StoreA", 2, 1, 1, 0, 3), ("StoreB", 1, 1, 1, 2, 4) ], ["Store", "MAX_A", "MAX_B", "MAX_C", "MAX_D", "N"])
需求分步实现
步骤1:移除门店已存在的ID
按分区(Div)过滤掉final_list中已在对应门店store_df里的ID:
# 构建门店-Div-ID的黑名单 store_blacklist = store_df.select("Store", "Div", "ID") # 过滤待选池:保留不在对应门店黑名单内的ID filtered_final_list = final_list.join( store_blacklist, on=["Div", "ID"], how="left_anti" ).withColumnRenamed("Rank", "Original_Rank")
步骤2:按Rank排序并按限额选ID
对每个门店,结合其分区的待选ID池,按原始排名顺序选取,同时控制单类别数量不超MAX、总数量不超N:
# 关联门店的分区信息与限额配置 store_div_max = store_df.select("Store", "Div").distinct().join( max_df, on="Store", how="inner" ) # 关联待选ID池与门店限额,形成每个门店的候选池 candidate_pool = store_div_max.join( filtered_final_list, on="Div", how="inner" ).orderBy("Original_Rank") # 定义窗口:按门店+类别分组算累计选量,按门店分组算总累计选量 cat_window = Window.partitionBy("Store", "Category").orderBy("Original_Rank") store_window = Window.partitionBy("Store").orderBy("Original_Rank") # 计算各类别及门店的累计选中数 candidate_with_counts = candidate_pool.withColumn( "Cat_Cum_Count", F.row_number().over(cat_window) ).withColumn( "Store_Cum_Count", F.row_number().over(store_window) ) # 筛选符合限额条件的ID selected_ids = candidate_with_counts.filter( (F.col("Cat_Cum_Count") <= F.when(F.col("Category") == "A", F.col("MAX_A")) .when(F.col("Category") == "B", F.col("MAX_B")) .when(F.col("Category") == "C", F.col("MAX_C")) .when(F.col("Category") == "D", F.col("MAX_D"))) & (F.col("Store_Cum_Count") <= F.col("N")) )
步骤3:更新max_df,新增实际分配列
统计每个门店各类别实际选中数量,合并到原限额表:
# 统计各类别实际分配数 actual_counts = selected_ids.groupBy("Store").agg( F.count(F.when(F.col("Category") == "A", 1)).alias("Add_A"), F.count(F.when(F.col("Category") == "B", 1)).alias("Add_B"), F.count(F.when(F.col("Category") == "C", 1)).alias("Add_C"), F.count(F.when(F.col("Category") == "D", 1)).alias("Add_D") ) # 合并到原限额表,空值填0 updated_max_df = max_df.join(actual_counts, on="Store", how="left").fillna(0, subset=["Add_A", "Add_B", "Add_C", "Add_D"])
步骤4:生成最终结果result_df
对每个门店的选中ID按原始排名排序,生成新排名:
result_window = Window.partitionBy("Store").orderBy("Original_Rank") result_df = selected_ids.select( "Store", "ID", "Category", F.row_number().over(result_window).alias("New_Rank") ).orderBy("Store", "New_Rank")
示例输出
updated_max_df(更新后的限额表)
| Store | MAX_A | MAX_B | MAX_C | MAX_D | N | Add_A | Add_B | Add_C | Add_D |
|---|---|---|---|---|---|---|---|---|---|
| StoreA | 2 | 1 | 1 | 0 | 3 | 1 | 1 | 1 | 0 |
| StoreB | 1 | 1 | 1 | 2 | 4 | 0 | 1 | 0 | 2 |
result_df(最终选中结果)
| Store | ID | Category | New_Rank |
|---|---|---|---|
| StoreA | ID2 | A | 1 |
| StoreA | ID4 | B | 2 |
| StoreA | ID6 | C | 3 |
| StoreB | ID4 | B | 1 |
| StoreB | ID7 | D | 2 |
| StoreB | ID8 | D | 3 |
内容的提问来源于stack exchange,提问作者Scope
相关产品推荐
相关产品推荐

