如何利用布尔值列表筛选PySpark DataFrame中的行
在PySpark中通过布尔列表筛选DataFrame行
PySpark是分布式计算框架,无法像Pandas那样直接用本地布尔列表索引筛选行——因为分布式DataFrame的行没有固定的物理顺序,本地列表的顺序无法直接对应到分布式数据的行。要实现需求,需要通过显式关联行号的方式将布尔列表与DataFrame的行一一绑定,具体步骤如下:
实现步骤
- 给原DataFrame添加连续的行号,确保每行有唯一标识;
- 将布尔列表转换为带对应行号的临时DataFrame;
- 关联两个DataFrame,筛选布尔值为
True的行,最后清理辅助列。
代码示例
1. 准备测试数据
from pyspark.sql import SparkSession from pyspark.sql.functions import row_number, monotonically_increasing_id, col from pyspark.sql.window import Window from pyspark.sql import Row # 初始化Spark会话 spark = SparkSession.builder.appName("BooleanFilterDemo").getOrCreate() # 原DataFrame示例 data = [("Geeks",), ("For",), ("Geeks",), ("is",), ("portal",), ("for",), ("Geeks",)] df1 = spark.createDataFrame(data, ["value"]) # 对应长度的布尔列表 unique_df1 = [True, False] * 3 + [True]
2. 给原DataFrame添加行号
这里用窗口函数生成连续行号,确保顺序与你生成布尔列表时的DataFrame顺序一致:
# 定义窗口:如果布尔列表是基于原DF的默认顺序,用monotonically_increasing_id()保证唯一排序 # 如果你是基于某列排序后生成的布尔列表,把orderBy的参数换成对应列(比如orderBy("value")) window = Window.orderBy(monotonically_increasing_id()) df_with_id = df1.withColumn("row_id", row_number().over(window))
3. 创建布尔列表的DataFrame
生成与原DF行号一一对应的布尔值DataFrame:
# 注意row_id从1开始(因为row_number()生成的行号从1起始) bool_data = [Row(row_id=i+1, keep=val) for i, val in enumerate(unique_df1)] bool_df = spark.createDataFrame(bool_data)
4. 关联并筛选
# 关联两个DataFrame,筛选需要保留的行,最后删除辅助列 new_df = df_with_id.join(bool_df, on="row_id", how="inner") \ .filter(col("keep") == True) \ .drop("row_id", "keep") # 查看结果 new_df.show()
运行后输出:
+-------+ | value| +-------+ | Geeks| | Geeks| | portal| | Geeks| +-------+
注意事项
- 行号的生成必须与你创建布尔列表时的DataFrame顺序完全匹配,如果布尔列表是基于原DF的某列排序结果生成的,一定要在窗口函数的
orderBy中指定该列,避免顺序错乱; - 如果原DataFrame数据量极大,
monotonically_increasing_id()可能比指定列排序更高效,但前提是你不需要严格的业务顺序,只需要保证列表与行的一一对应。
内容的提问来源于stack exchange,提问作者Clark Sims
相关产品推荐
相关产品推荐

