非SQL方式:基于另一DataFrame过滤PySpark DataFrame
不用SQL过滤PySpark DataFrame的解决方案
没问题,咱们纯用PySpark DataFrame API就能实现这个需求,完全不需要写SQL语句。我给你一步步拆解实现过程:
步骤1:初始化Spark环境并创建测试数据
首先咱们先把你提供的df1和df2数据转换成PySpark DataFrame,方便后续操作:
from pyspark.sql import SparkSession # 初始化SparkSession(PySpark操作的核心入口) spark = SparkSession.builder.appName("FilterDFWithoutSQL").getOrCreate() # 构建df1的数据源 df1_data = [ (1, "2020-12-01", None, "Paul"), (2, "2020-12-02", None, "Mary"), (3, "2020-12-03", None, "David"), (4, "2020-12-04", None, "Marley") ] df1 = spark.createDataFrame(df1_data, schema=["id", "create", "change", "name"]) # 构建df2的数据源 df2_data = [ (1, "2020-12-01", "2020-12-30", "Paul"), (2, "2020-12-02", None, "Mary"), (3, "2020-12-03", None, "David"), (4, "2020-12-04", "2020-12-30", "Marley"), (5, "2020-12-30", None, "Ted") ] df2 = spark.createDataFrame(df2_data, schema=["id", "create", "change", "name"])
步骤2:实现过滤逻辑
咱们的目标是排除df2中同时满足两个条件的行:
change字段值为2020-12-30id存在于df1的id集合中
这里我提供两种高效的实现方式,你可以根据数据量选择:
方式一:直接用isin+过滤条件
这种方式代码直观,适合中小数据量场景:
# 提取df1中所有的id(转换成Python列表) df1_id_list = [row.id for row in df1.select("id").collect()] # 构建过滤条件:取反(排除符合两个条件的行) df3 = df2.filter( ~((df2.change == "2020-12-30") & (df2.id.isin(df1_id_list))) )
方式二:用广播变量+左反连接(大数据量推荐)
如果你的数据量很大,用广播变量+左反连接能大幅提升性能,避免重复传输数据:
from pyspark.sql.functions import broadcast # 先筛选出df2中需要排除的候选行(change=2020-12-30),再和df1关联得到要排除的行 exclude_rows = df2.filter(df2.change == "2020-12-30").join(broadcast(df1), on="id", how="inner") # 用左反连接从df2中排除掉这些行 df3 = df2.join(exclude_rows, on=["id", "create", "change", "name"], how="left_anti")
步骤3:查看结果
执行下面的代码就能看到最终的df3数据:
df3.show()
输出结果(符合预期,保留了无需排除的行):
+---+----------+------+-----+ | id| create|change| name| +---+----------+------+-----+ | 2|2020-12-02| null| Mary| | 3|2020-12-03| null|David| | 5|2020-12-30| null| Ted|
(注:你提供的目标示例里漏了Ted的行,按照规则他的行不符合排除条件,应该被保留哦)
内容的提问来源于stack exchange,提问作者Paulo
相关产品推荐
相关产品推荐

