PySpark内存列表执行filter+lambda时抛出NOT_COLUMN_OR_STR异常排查
问题分析与解决:collect()后仍触发PySpark异常
为什么会触发PySpark异常?
你以为collect()得到的是纯Python原生类型的列表,但实际上currentCustomers里的元素是PySpark的Row对象,不是普通的字符串或数值。当你在lambda里写y not in x时,Row对象的in操作被Spark错误识别成了SQL表达式逻辑,所以抛出NOT_COLUMN_OR_STR异常——Spark误以为你在操作Column对象,而非Python原生对象。
collect()确实会加载数据到内存
collect()的作用就是把DataFrame的所有数据拉取到Driver节点的内存中,转换成Python列表,但默认情况下,列表里的元素是PySpark Row对象,不是原生数据类型。比如你取customer_id列,collect()得到的是[Row(customer_id='1001'), Row(customer_id='1002')],而不是['1001', '1002']。
解决方法
方法1:修改collectListItems,直接提取原生类型
在收集列表时,直接把Row里的列值取出来,得到纯Python列表:
def collectListItems(df, column_name): # 提取指定列的原生值,返回纯Python列表 return df.select(column_name).rdd.flatMap(lambda row: row).collect()
这样currentCustomers和activeCustomers就是纯字符串/整数列表,执行filter时不会触发PySpark异常。
方法2:在filter时从Row中提取值
如果不想修改收集方法,就在过滤时手动取出Row里的具体值:
# 假设列名为customer_id,用属性方式取值 list(filter(lambda x: x.customer_id not in activeCustomers, currentCustomers)) # 或者用索引(如果确定是第一列) list(filter(lambda x: x[0] not in activeCustomers, currentCustomers))
内容的提问来源于stack exchange,提问作者Shane McGarry
相关产品推荐
相关产品推荐

