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

PySpark实现:判断DataFrame列值是否存在于另一DataFrame列中

PySpark:标记DataFrame中值是否存在于另一个DataFrame的实现方法

数据初始化代码

先给出可复现的DataFrame创建代码,方便快速验证:

from pyspark.sql import SparkSession
from pyspark.sql.functions import col, when, broadcast

# 初始化Spark会话
spark = SparkSession.builder.appName("FlagValueCheck").getOrCreate()

# 创建df1
data1 = [("the",), ("quick",), ("brown",), ("fox",), ("brown",)]
df1 = spark.createDataFrame(data1, ["value"])

# 创建df2
data2 = [("brown",), ("fox",), ("brown",)]
df2 = spark.createDataFrame(data2, ["value"])

方法一:广播关联(推荐,高效适配多数场景)

当df2数据量较小时,用广播关联+左连接的方式,既能保证效率,又能避免Driver节点内存溢出:

# 先对df2去重,减少关联数据量
df2_unique = df2.select("value").distinct()

# 广播df2_unique后与df1左连接,生成flag列
result = df1.join(broadcast(df2_unique), on="value", how="left") \
            .withColumn("flag", when(col("value").isNotNull(), 1).otherwise(0))

# 查看结果
result.show()

方法二:isin函数(仅适用于df2数据量极小的场景)

如果df2数据量非常小,可以直接把值拉取到Driver节点用isin判断,但数据量大时容易引发内存问题,谨慎使用:

# 获取df2的唯一值列表
target_values = [row.value for row in df2.select("value").distinct().collect()]

# 生成flag列
result = df1.withColumn("flag", when(col("value").isin(target_values), 1).otherwise(0))

result.show()

最终输出结果

两种方法都会得到期望的输出:

+-----+----+
|value|flag|
+-----+----+
|  the|   0|
|quick|   0|
|brown|   1|
|  fox|   1|
|brown|   1|
+-----+----+

内容的提问来源于stack exchange,提问作者ZZZZZZZZZ

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 02:27:12