PySpark如何判断DataFrame某列是否包含指定值?
解决PySpark DataFrame判断列是否包含指定值的问题
我来帮你搞定这个问题!你遇到的核心问题是PySpark的filter()方法返回的永远是一个DataFrame对象——哪怕过滤后没有任何数据,这个对象在Python的布尔判断里也会被视为True,所以不管有没有匹配到"3",你的代码都会输出"Yes"。
下面给你几种正确的实现方式,按效率和简洁性排序:
方法1:用exists()高效判断(推荐)
exists()方法会检查DataFrame是否存在至少一行数据,而且它会在找到第一个匹配项后就停止计算,非常适合大数据场景:
if df.where(df.id == "3").exists(): print('Yes') else: print('No')
或者用col()函数写法更规范一点:
from pyspark.sql.functions import col if df.where(col("id") == "3").exists(): print('Yes') else: print('No')
方法2:用count()统计匹配行数
如果需要知道具体有多少匹配行,或者你习惯这种写法,可以用count()判断是否大于0:
if df.filter(df.id == "3").count() > 0: print('Yes') else: print('No')
注意:count()会遍历所有匹配的行,数据量很大的时候效率不如exists()。
方法3:用isin()处理多值判断
如果需要判断值是否在一个集合里(比如同时检查"3"、"4"),isin()就很有用,用法是传入一个列表,再结合exists()或count():
# 检查id是否包含"3" if df.where(df.id.isin(["3"])).exists(): print('Yes') else: print('No') # 扩展:检查id是否包含"3"或"4" if df.where(df.id.isin(["3", "4"])).exists(): print('Yes') else: print('No')
为什么Pandas的写法不能直接照搬?
Pandas是单机处理,'3' in df_init['id'].values可以直接在本地内存里检查值是否存在;但PySpark是分布式计算,所有操作都要通过Spark的API生成执行计划,然后在集群上运行,不能直接用Python的in操作符来判断分布式数据。
内容的提问来源于stack exchange,提问作者George
相关产品推荐
相关产品推荐

