如何对PySpark DataFrame执行过滤操作且保留DataFrame格式
PySpark过滤操作后保留DataFrame格式的解决方法
问题原因
你执行代码后得到列表是因为末尾调用了collect()方法:该方法属于PySpark的行动(Action)算子,作用是将分布式存储的DataFrame全量数据拉取到当前驱动节点的内存中,返回值本身就是由Row对象组成的Python列表,因此会覆盖原来的DataFrame类型变量。
解决方式
直接删除代码末尾的.collect()调用即可,调整后代码如下:
datalabel = datalabel.filter(datalabel.subs_no.isNotNull())
调整后赋值得到的datalabel仍然是标准PySpark DataFrame类型,可正常执行后续的转换、写入等DataFrame操作。
补充建议
如果需要验证过滤结果是否符合预期,不要直接使用collect()拉取全量数据,推荐搭配limit()算子取少量数据预览:
# 仅预览前10行结果,不会修改原datalabel的DataFrame属性 datalabel.limit(10).show()
该方式既可以快速校验过滤逻辑,也能避免全量数据拉取可能导致的驱动节点内存溢出问题。
内容的提问来源于stack exchange,提问作者Nabih Bawazir
相关产品推荐
相关产品推荐

