PySpark 3.0.1与3.1.2版本filter算子性能异常问题咨询
问题根因
你的代码实现没有语法问题,该故障是Spark 3.1.x版本的已知内核bug,和业务逻辑无关。
Spark 3.1.0到3.1.2版本的Catalyst优化器在处理字符串前缀匹配类算子时,会出现逻辑计划递归转换死循环,驱动端内存被持续占满后就会抛出GC overhead limit exceeded错误,你调用explain()触发执行计划生成时刚好命中了这个递归逻辑。
修复方案
按优先级可选择以下方案:
- 优先替换字符串筛选逻辑,改用
startswith方法实现前缀匹配,该写法不会触发3.1.x版本的优化器bug,且执行效率和你原先的like('Invalid%')完全一致:
df_filtered = df.filter(col('val').startswith('Invalid'))
- 如果需要保留原有写法,可以在SparkSession初始化时添加配置,临时关闭触发bug的优化规则,该方案会小幅影响整体SQL优化性能,仅建议临时过渡使用:
from pyspark.sql import SparkSession spark = SparkSession.builder \ .config("spark.sql.optimizer.excludedRules", "org.apache.spark.sql.catalyst.optimizer.OptimizeIn") \ .getOrCreate()
- 长期使用建议升级Spark版本到3.2.0及以上,该bug已在3.2.0版本被官方修复,升级后原有写法可正常运行。
内容的提问来源于stack exchange,提问作者Kashyapgv
相关产品推荐
相关产品推荐

