PySpark DataFrame用.where/.filter过滤速度慢,有无更高效的处理方案?
PySpark过滤计算列耗时过高优化方案
核心问题根因
你观察到的过滤步骤耗时3小时并非过滤操作本身导致,本质是Spark惰性计算机制造成的感知误差:之前3分钟跑完的RapidFuzz模糊匹配仅完成了逻辑执行计划生成,并未真正触发全量数据计算,直到调用where/filter执行过滤时,才会从头到尾执行从数据源读取、模糊匹配计算到过滤输出的全链路流程,相当于把模糊匹配的计算成本全部叠加到了过滤步骤上。
优化方案
- 提前缓存模糊匹配中间结果
包含计算列l_levensthien的l1是中间表,在做过滤前先将l1缓存到内存/磁盘,避免每次执行过滤操作都重复运行RapidFuzz计算逻辑,代码示例:from pyspark import StorageLevel # 生成l1后立即缓存,内存不足时选择MEMORY_AND_DISK存储级别 l1 = l1.persist(StorageLevel.MEMORY_AND_DISK) # 触发轻量计算完成缓存落地,耗时和你之前统计的模糊匹配耗时一致(约3分钟) l1.count() # 此时再执行过滤,仅需数秒到数分钟即可完成 l2 = l1.where(l1.l_levensthien != 0) # l1不再使用时释放缓存 # l1.unpersist() - 前置过滤逻辑实现谓词下推
若存在源数据层面的过滤条件,比如空值过滤、无效行剔除,在计算l_levensthien列之前就完成过滤,减少模糊匹配的计算总量。 - 优化计算列数据类型
确认l_levensthien列的类型为整型,整型的数值比较运算效率远高于字符串、浮点型等其他数据类型,可大幅降低过滤操作的计算开销。 - 调整并行度配置
检查l1的分区数量,保持分区数为集群总CPU核心数的2-3倍,避免分区数过少导致并行度不足、分区数过多导致调度开销过高的问题。
内容的提问来源于stack exchange,提问作者moe94z
相关产品推荐
相关产品推荐

