基于分组+过滤条件生成Spark DataFrame布尔列
解决方案
要生成target_is_correct布尔列,核心是找到每行i-j差值对应的同组目标行id,再与当前行target_id对比。以下提供两种Spark原生实现方法:
方法一:自连接实现
适合大数据量场景,利用Spark分布式连接能力处理:
from pyspark.sql import functions as F # 1. 计算每行的i-j差值 df_with_diff = df.withColumn("diff", F.col("i") - F.col("j")) # 2. 创建group-i-id的映射表,重命名列避免冲突 id_mapping = df.select("group", "i", "id") \ .withColumnRenamed("i", "target_i") \ .withColumnRenamed("id", "correct_target_id") # 3. 自连接匹配同组中i等于diff的目标行,对比生成布尔列 result_df = df_with_diff.join( id_mapping, (df_with_diff["group"] == id_mapping["group"]) & (df_with_diff["diff"] == id_mapping["target_i"]), "left" ).withColumn("target_is_correct", F.col("target_id") == F.col("correct_target_id")) \ .drop("diff", "target_i", "correct_target_id", id_mapping["group"]) # 查看结果 result_df.show()
代码说明
- 先计算
diff列存储i-j的结果; - 映射表
id_mapping存储每个分组下i值对应的正确id; - 通过自连接关联同组且
diff等于target_i的行,直接对比target_id和correct_target_id得到布尔值; - 最后清理中间生成的冗余列。
方法二:窗口函数+Map映射
代码更简洁,利用窗口函数在分组内构建映射关系:
from pyspark.sql import functions as F from pyspark.sql.window import Window # 定义窗口:按group分组 group_window = Window.partitionBy("group") # 1. 每个分组内构建i到id的映射字典 df_with_map = df.withColumn( "id_map", F.map_from_entries(F.collect_list(F.struct("i", "id")).over(group_window)) ) # 2. 计算diff,从映射字典中取出正确id并对比 result_df = df_with_map.withColumn("diff", F.col("i") - F.col("j")) \ .withColumn("correct_target_id", F.col("id_map")[F.col("diff")]) \ .withColumn("target_is_correct", F.col("target_id") == F.col("correct_target_id")) \ .drop("id_map", "diff") # 查看结果 result_df.show()
代码说明
- 窗口函数
collect_list在每个分组内收集所有(i, id)结构,通过map_from_entries转成键值对映射; - 计算
diff后,直接通过字典索引取出对应的正确id; - 对比
target_id和正确id生成布尔列,最后清理中间列。
两种方法运行后都会得到预期结果:
+-----+-----+---+---+---------+-----------------+ | id|group| i| j|target_id|target_is_correct| +-----+-----+---+---+---------+-----------------+ |id-A5| A| 5| 0| id-A5| true| |id-A7| A| 7| 2| id-A5| true| |id-B0| B| 0| 0| id-B0| true| |id-B1| B| 1| 1| id-B0| true| |id-B2| B| 2| 1| id-B0| false| +-----+-----+---+---+---------+-----------------+
内容的提问来源于stack exchange,提问作者L.B.
相关产品推荐
相关产品推荐

