PySpark如何将N行分配到X组且相同B+C组合不拆分并分配日期D
实现方案
核心逻辑
你之前的代码拆分了同B+C组合的原因是直接对每行生成独立行号,同一组合的不同行会拿到不同行号,取模后自然分配到不同组。要满足约束,需先为所有唯一的(B,C)组合统一分配组号,再关联回原数据,确保同一组合的所有行拿到相同组号。
具体实现代码
from pyspark.sql import functions as F from pyspark.sql.window import Window # 上游定义的常量参数X X = 3 # 步骤1:生成唯一(B,C)组合与组号的映射表 bc_group_map = df1.select("B", "C").distinct() \ # 对BC组合随机排序,满足要求中组合排序可随机的规则 .withColumn("bc_seq", F.row_number().over(Window.orderBy(F.rand()))) \ # 组号范围为1~X,刚好对应日期偏移量 .withColumn("group_id", (F.col("bc_seq") % X) + 1) \ .drop("bc_seq") # 步骤2:将组号映射回原数据 df2 = df1.join(bc_group_map, on=["B", "C"], how="left") # 步骤3:生成日期列D,两种格式按需选择 # 格式1:真实日期,当日+group_id天 # df2 = df2.withColumn("D", F.date_add(F.current_date(), F.col("group_id"))) # 格式2:示例中的date1/date2字符串格式 df2 = df2.withColumn("D", F.concat(F.lit("date"), F.col("group_id"))) # 可选:删除不需要的中间列 df2 = df2.drop("group_id") # 查看结果 df2.show()
均匀性优化(可选)
如果不同(B,C)组合的行数差异较大,可先统计每个组合的行数,按行数降序排序后再分配组号,避免大组集中在同一个分组,进一步提升各分组大小的均匀度:
bc_group_map = df1.groupBy("B", "C").agg(F.count("*").alias("row_cnt")) \ .orderBy(F.desc("row_cnt")) \ .withColumn("bc_seq", F.monotonically_increasing_id()) \ .withColumn("group_id", (F.col("bc_seq") % X) + 1) \ .drop("row_cnt", "bc_seq")
效果说明
- 同一(B,C)组合的所有行永远分配到同一个组,不会拆分,满足约束
- 当X大于唯一(B,C)组合的总数量时,多余的分组会自动空置,符合你给出的X=4、X=5的示例场景
- 分组分配完全随机,可通过调整
orderBy的规则自定义排序逻辑
内容的提问来源于stack exchange,提问作者ionah
相关产品推荐
相关产品推荐

