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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.27 03:45:03