You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.22 09:38:46