Pandas转PySpark代码求助:筛选指定student_id的DataFrame
PySpark实现Pandas筛选逻辑的方案
先明确原Pandas代码的逻辑:
- 筛选
house列值为True的行,提取其中的student_id并去重 - 从原DataFrame中筛选所有
student_id属于上述去重集合的行
下面提供三种PySpark实现方式,适配不同数据量场景:
方法一:基于collect+isin(小数据量适用)
直接复刻Pandas的思路,先提取符合条件的student_id集合,再进行筛选:
# 获取house为True的唯一student_id集合 uni_id = df.filter(df.house == True).select("student_id").distinct().rdd.flatMap(lambda x: x).collect() # 筛选目标行 class_df = df.filter(df.student_id.isin(uni_id))
注意:collect()会把数据拉取到Driver节点,数据量大时可能引发内存问题,仅适合小数据集。
方法二:基于Join(大数据量推荐)
通过Spark分布式join操作实现,避免数据拉取到Driver:
# 生成包含目标student_id的临时DataFrame target_students = df.filter(df.house == True).select("student_id").distinct() # 内连接筛选出符合条件的行 class_df = df.join(target_students, on="student_id", how="inner")
方法三:基于窗口函数(更灵活的分布式方案)
通过窗口函数标记每个student_id组是否存在house为True的记录,再筛选:
from pyspark.sql import Window import pyspark.sql.functions as F # 按student_id分组的窗口 student_window = Window.partitionBy("student_id") # 标记组内是否有house为True的记录,再筛选 class_df = df.withColumn("has_house_true", F.max(F.col("house")).over(student_window)) \ .filter(F.col("has_house_true") == True) \ .drop("has_house_true")
这种方式全程在分布式节点处理,性能最优,适合大规模数据集。
内容的提问来源于stack exchange,提问作者Abhinav bharti
相关产品推荐
相关产品推荐

