PySpark如何按多条件匹配两个DataFrame筛选符合要求的行
PySpark 筛选匹配行实现方案
有两种常用的实现方式,可根据你的数据量场景选择:
方案1:内连接实现(推荐,适配所有数据量场景)
直接以三个需要匹配的列作为关联键做内连接,再仅保留d1的原字段即可,Spark对join算子有底层优化,大数据量下性能表现更好。
示例代码:# 声明匹配用的公共字段 match_keys = ["c1", "c2", "value"] # 内连接后筛选d1的所有列,过滤掉匹配失败的行 filtered_d1 = d1.join(d2, on=match_keys, how="inner").select(d1["*"])方案2:小表场景下的isin过滤实现
如果d2的数据量很小可以完全加载到driver内存,也可以用collect后匹配的方式,不过大数据量下不推荐,容易出现内存溢出。
示例代码:from pyspark.sql import functions as F # 把d2的匹配列拼接为数组结构,collect后广播 d2_match_list = d2.select(F.array("c1", "c2", "value")).rdd.flatMap(lambda x: x).collect() # 过滤d1中匹配列拼接后在列表内的行 filtered_d1 = d1.filter(F.array("c1", "c2", "value").isin(d2_match_list))
注意事项
- 如果你的业务要求空值也判定为相等(比如d1的c1为null、d2的c1为null也视为匹配),需要使用
eqNullSafe方法改写关联条件,原生的相等判断遇到null会返回false:from pyspark.sql import functions as F filtered_d1 = d1.join( d2, (F.eqNullSafe(d1.c1, d2.c1)) & (F.eqNullSafe(d1.c2, d2.c2)) & (F.eqNullSafe(d1.value, d2.value)), how="inner" ).select(d1["*"]) - 如果d2存在重复的匹配行,join后会导致d1的行重复,如需去重可在select后追加
.dropDuplicates()即可。
内容的提问来源于stack exchange,提问作者MarsEclipse
相关产品推荐
相关产品推荐

