Pyspark如何无需转Pandas快速删除DataFrame中单词数不足n的行
实现方案
你原方案效率极低的核心原因是:将Spark分布式DataFrame转为Pandas DataFrame时,会把全量600万行数据全部拉取到Driver节点单节点处理,既容易触发Driver内存溢出,单线程遍历的效率也远低于分布式计算。
直接使用PySpark原生函数即可实现分布式过滤,无需转换数据格式,性能提升非常明显:
from pyspark.sql.functions import split, size, trim, col # 过滤text字段单词数≥4的行,同时自动过滤空text、全空格text的无效记录 cleaned_df = df.filter( size(split(trim(col("text")), "\\s+")) >= 4 )
逻辑说明
trim(col("text")):先去除文本首尾的空白字符,避免首尾空格分割出空字符串干扰单词计数split(..., "\\s+"):按1个或多个空白字符分割文本,避免文本内连续空格产生的空字符串被计为有效单词size(...):计算分割后的单词数组长度,直接做过滤条件
所有计算都在集群Executor节点并行执行,不需要将数据拉取到Driver端,处理600万行数据的耗时仅为原方案的几十分之一。
如果需要额外显式过滤text为null的记录,可以补充判断条件:
cleaned_df = df.filter( col("text").isNotNull() & (size(split(trim(col("text")), "\\s+")) >= 4) )
内容的提问来源于stack exchange,提问作者gigioneggiavamo
相关产品推荐
相关产品推荐

