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
相关产品推荐
相关产品推荐

