PySpark实现:筛选存在跨行双条件的DataFrame ID对应行
解决方法:筛选同时存在
NEW=1和LEFT=1的ID的所有行 针对你的需求,我们可以通过高效的方式定位目标ID,再提取对应所有记录——考虑到你的数据集超过1300万条,我们优先选择性能友好的实现方案,避免不必要的计算开销。
方法一:分组聚合筛选ID(推荐,适合大数据量)
这种方式先通过分组聚合锁定符合条件的ID,再与原表关联,在大数据场景下性能更优:
from pyspark.sql import SparkSession from pyspark.sql.functions import max, col # 初始化SparkSession(若未初始化) spark = SparkSession.builder.appName("FilterQualifiedIDs").getOrCreate() # 假设你的原始DataFrame名为df # 第一步:找出同时存在NEW=1和LEFT=1的ID qualified_ids = df.groupBy("ID") .agg( max(col("NEW")).alias("has_new"), max(col("LEFT")).alias("has_left") ) .filter((col("has_new") == 1) & (col("has_left") == 1)) .select("ID") # 第二步:关联原表,获取这些ID的所有行 result = df.join(qualified_ids, on="ID", how="inner") # 查看结果 result.show()
思路说明:
groupBy("ID")按ID分组后,用max(NEW)和max(LEFT)判断该ID是否存在至少一行NEW=1或LEFT=1(只要有一行满足,最大值就是1)。- 过滤出同时满足两个条件的ID,再通过
inner join关联原表,这种关联操作在Spark中针对大数据量做了深度优化,效率很高。
方法二:窗口函数实现
如果你更习惯用窗口函数,也可以通过为每行标记所属ID的状态来过滤:
from pyspark.sql.window import Window from pyspark.sql.functions import max, col # 定义窗口:按ID分区 window_spec = Window.partitionBy("ID") # 为每行添加标记,指示该ID是否有NEW=1和LEFT=1的记录 df_with_flags = df.withColumn("has_new", max(col("NEW")).over(window_spec)) .withColumn("has_left", max(col("LEFT")).over(window_spec)) # 过滤出符合条件的行,再移除临时标记列 result = df_with_flags.filter((col("has_new") == 1) & (col("has_left") == 1)) .drop("has_new", "has_left") result.show()
思路说明:
- 窗口函数会为每一行计算其所属ID的
NEW和LEFT最大值,让每行都能知道自己的ID是否符合要求。 - 过滤后移除临时标记列即可得到结果,这种方式无需额外关联,但在极端大数据量下,性能略逊于分组聚合方案。
性能优化提示
针对1300万条记录的数据集,建议:
- 确保Spark集群配置足够的executor内存和核心数,避免资源瓶颈。
- 如果ID列重复率较高,优先选择分组聚合方案,它会先对ID去重再关联,节省计算资源。
- 可提前对ID列进行分区或分桶,进一步提升关联/窗口计算的效率。
最终你会得到预期结果:
+---+----------+------+----+---+ | ID| DATE|ACTIVE|LEFT|NEW| +---+----------+------+----+---+ |456|2021-03-01| 1| 0| 1| |456|2021-06-01| 1| 1| 0| +---+----------+------+----+---+
内容的提问来源于stack exchange,提问作者Alf
相关产品推荐
相关产品推荐

