PySpark中基于另一DataFrame列值存在性新增列的实现问题
解决PySpark判断列值是否存在于另一个DataFrame列中的问题
我来帮你搞定这个问题!你之前的思路是对的,但在PySpark里直接用expr引用另一个DataFrame的列是行不通的——因为SQL表达式的上下文没法直接访问其他DF的列集合。下面给你几种靠谱的解决方案,适配不同的场景:
方案一:提取目标集合后用when/otherwise(适合小数据集)
如果bears2的week_if数据量不大,可以先把这些值提取成Python集合,再用isin配合when函数判断:
import pyspark.sql.functions as F # 提取bears2中去重的week_if值,转成Python列表 bears_week_ids = [row.week_if for row in bears2.select("week_if").distinct().collect()] # 生成新列:匹配到则为1,否则为0 df3 = df2.withColumn( "game", F.when(F.col("week_id").isin(bears_week_ids), 1).otherwise(0) )
⚠️ 注意:如果bears2数据量很大,collect()会把数据拉到Driver节点,可能导致内存溢出,这种情况建议用下面的join方案。
方案二:左连接+coalesce(适合大数据集)
通过左连接关联两个DataFrame,再用coalesce填充未匹配的情况,这是Spark处理大数据场景的标准做法:
import pyspark.sql.functions as F # 先对bears2去重,提前标记匹配值为1 bears_unique = bears2.select("week_if").distinct().withColumn("game", F.lit(1)) # 左连接df2和处理后的bears_unique,填充未匹配的0 df3 = df2.join( bears_unique, df2.week_id == bears_unique.week_if, how="left" ).withColumn( "game", F.coalesce(F.col("game"), F.lit(0)) # 把null转成0 ).drop("week_if") # 删掉多余的连接列
如果bears_unique数据量很小,可以用broadcast优化连接性能:
df3 = df2.join( F.broadcast(bears_unique), df2.week_id == bears_unique.week_if, how="left" )... # 后续步骤同上
方案三:用expr实现(需先提取值)
如果你坚持想用expr,可以先把week_if的值拼接成SQL能识别的字符串,再传入表达式:
import pyspark.sql.functions as F # 提取值并转成逗号分隔的字符串 bears_week_str = ",".join(map(str, [row.week_if for row in bears2.select("week_if").distinct().collect()])) # 用expr生成新列 df3 = df2.withColumn( "game", F.expr(f"IF(week_id IN ({bears_week_str}), 1, 0)") )
这个方案和方案一类似,同样只适合小数据集。
为什么你之前的代码失败?
你之前写的F.expr("""IF(df2.week_id IN(bears2.week_if), 1, 0))""")无法生效,是因为Spark SQL表达式只能访问当前DataFrame的列,没法直接引用另一个DataFrame的列对象。必须先把另一个DF的目标列值提取出来,或者通过join的方式把两个DF的列放到同一个上下文里。
内容的提问来源于stack exchange,提问作者mlewis
相关产品推荐
相关产品推荐

