PySpark使用isin判断列值是否存在于同行列表列的问题问询
PySpark判断单列值是否在同行数组列的解决方案
报错原因
你的写法错误核心有两点:
isin()方法接收的是固定字面量列表,不支持直接传入DataFrame的数组列作为判断依据collect()是Action操作,会把整列数据拉取到Driver节点生成Python列表,无法用于逐行判断逻辑,直接写在列操作代码中本身属于语法错误
解决方案
使用PySpark内置的array_contains函数即可实现需求,该函数专门用于逐行判断某个值是否存在于当前行的数组列中。
完整修正代码
首先导入依赖:
from pyspark.sql import functions as F from pyspark.sql.window import Window
业务逻辑代码:
# 注意修正原窗口范围的拼写错误,原写法unboundedPreceeding拼写错误,正确为unboundedPreceding w = Window.partitionBy('ID').orderBy('date').rowsBetween(Window.unboundedPreceding, -1) df = df.withColumn('main_list', F.collect_set('loc').over(w)) \ .withColumn('GOAL_f', F.array_contains(F.col('main_list'), F.col('loc')).cast('int'))
array_contains会返回布尔值结果,使用cast('int')可直接转为你示例要求的1/0格式,输出结果和你给出的样例完全一致。
内容的提问来源于stack exchange,提问作者s223
相关产品推荐
相关产品推荐

