PySpark:如何用RDD.filter()实现WHERE条件过滤操作?
用Spark RDD的filter()实现WHERE条件过滤(Python版)
嘿,我懂你现在的困惑——想用Spark RDD的filter()来实现类似SQL里WHERE的过滤逻辑,却不知道怎么下手对吧?结合你给的CSV数据,我一步步给你演示怎么写:
首先,咱们得先把CSV数据正确加载成RDD,还要跳过表头(不然表头会被当成数据行处理,转数字的时候会报错)。假设你的CSV文件路径是your_data.csv,先写加载部分:
from pyspark import SparkContext # 初始化SparkContext sc = SparkContext("local", "RDDFilterDemo") # 加载CSV文件 raw_rdd = sc.textFile("your_data.csv") # 提取表头并过滤掉表头行 header = raw_rdd.first() data_rdd = raw_rdd.filter(lambda line: line != header) \ .map(lambda line: [field.strip() for field in line.split(",")])
这里用map()把每行分割成字段列表,还顺便用strip()去掉了字段前后的空格(看你的CSV里有, Yves,这种带空格的格式,处理后更方便后续判断)。
接下来就是核心的filter()用法了,它接受一个函数(通常用lambda),返回True的元素会被保留,完全对应SQL的WHERE条件:
1. 单条件过滤(比如筛选年龄大于30的记录)
# 筛选age>30的记录,注意要把字符串转成整数 age_over_30_rdd = data_rdd.filter(lambda row: int(row[3]) > 30)
这里row[3]对应CSV里的age列(索引从0开始:id是0,firstname是1,city是2,age是3...),转成整数才能做数值比较。
2. 多条件过滤(比如职业是"Pilote de chasse"且薪资>2000)
# 多条件用and/or连接,注意字符串匹配要准确 target_rdd = data_rdd.filter( lambda row: row[4] == "Pilote de chasse" and int(row[5]) > 2000 )
这里row[4]是job列,row[5]是salary列,两个条件同时满足才会被保留。
3. 模糊匹配(比如筛选城市名称包含"OLS"的记录)
如果需要类似SQL里的LIKE,直接用字符串的in或者find()就行:
city_contains_ols_rdd = data_rdd.filter(lambda row: "OLS" in row[2])
最后,咱们可以把过滤后的结果收集打印出来看看:
# 打印年龄大于30的结果 print("年龄大于30的人员记录:") for record in age_over_30_rdd.collect(): print(record) # 打印多条件过滤的结果 print("\n符合职业和薪资要求的记录:") for record in target_rdd.collect(): print(record)
需要注意的几个点:
- 字段索引一定要对应正确,别搞混列的顺序;
- 涉及数值比较的字段,必须先转成对应的类型(int/float),不然会报错;
- 如果CSV里有缺失值,最好在过滤前加个判断,比如
if row[3].isdigit()避免转类型失败。
内容的提问来源于stack exchange,提问作者Essex
相关产品推荐
相关产品推荐

