Spark按条件随机选行:分组按颜色优先级选行并限制指定颜色数量
解决方案
核心实现逻辑
- 首先给颜色映射优先级权重:
blue:1 > green:2 > yellow:3 > red:4,数值越小优先级越高 - 单独处理A分组的blue限制:对A分组下的所有blue行先做随机排序,仅保留最多2行的最高优先级,剩余的blue行优先级调整为低于red,避免被优先选取
- 最后按
prod_name分组,窗口内按调整后的优先级升序、同优先级随机排序,取每个分组前3行即可
完整PySpark实现代码
from pyspark.sql import SparkSession from pyspark.sql.functions import col, when, rand, row_number from pyspark.sql.window import Window # 初始化SparkSession spark = SparkSession.builder.appName("group_select_with_rule").getOrCreate() # 构造示例DataFrame data = [ ("A", "blue", 100), ("A", "blue", 200), ("A", "blue", 300), ("A", "blue", 300), ("A", "yellow", 309), ("B", "green", 408), ("B", "blue", 50), ("C", "red", 6000), ("C", "blue", 70), ("C", "green", 10) ] df = spark.createDataFrame(data, schema=["prod_name", "colour", "prod_id"]) # 第一步:映射原始颜色优先级 colour_priority = when(col("colour") == "blue", 1)\ .when(col("colour") == "green", 2)\ .when(col("colour") == "yellow", 3)\ .when(col("colour") == "red", 4)\ .otherwise(99) df = df.withColumn("original_priority", colour_priority) # 第二步:处理A分组的blue数量限制,调整优先级 a_blue_window = Window.partitionBy("prod_name", "colour").orderBy(rand()) df = df.withColumn("a_blue_rn", when( (col("prod_name") == "A") & (col("colour") == "blue"), row_number().over(a_blue_window) ).otherwise(0) ) # A分组中序号超过2的blue行优先级调整为5,低于最低的red优先级 df = df.withColumn("adjusted_priority", when( (col("prod_name") == "A") & (col("colour") == "blue") & (col("a_blue_rn") > 2), 5 ).otherwise(col("original_priority")) ) # 第三步:分组内按调整后优先级、随机排序取前3行 final_window = Window.partitionBy("prod_name").orderBy("adjusted_priority", rand()) df = df.withColumn("final_rn", row_number().over(final_window)) result_df = df.filter(col("final_rn") <= 3).select("prod_name", "colour", "prod_id") # 输出结果 result_df.show()
结果验证
针对提供的示例数据,执行后完全符合所有规则:
- A分组最多选2个随机blue行,剩下1个名额选yellow行,共3行
- B分组优先选blue行,再选green行,凑够3行(数据不足时按实际数量返回)
- C分组优先选blue、green行,再选red行,凑够3行
内容的提问来源于stack exchange,提问作者user3735871
相关产品推荐
相关产品推荐

