You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

PySpark 3.2中slice函数触发SparkRuntimeException问题排查求助

PySpark 3.2中slice操作触发SparkRuntimeException的原因分析与排查方向

问题重现

使用PySpark 3.2版本对DataFrame中的数组执行slice操作时,触发以下异常:

org.apache.spark.SparkRuntimeException: Unexpected value for length in function slice: length must be greater than or equal to 0

但逻辑上已经通过df.filter(F.col('end_index') > 0)确保传入slice的length值为正。

完整异常栈

Caused by: org.apache.spark.SparkRuntimeException: Unexpected value for length in function slice: length must be greater than or equal to 0.
                at org.apache.spark.sql.errors.QueryExecutionErrors$.unexpectedValueForLengthInFunctionError(QueryExecutionErrors.scala:1602)
                at org.apache.spark.sql.errors.QueryExecutionErrors.unexpectedValueForLengthInFunctionError(QueryExecutionErrors.scala)
                at org.apache.spark.sql.catalyst.expressions.GeneratedClass$SpecificPredicate.Slice_0$(Unknown Source)
                at org.apache.spark.sql.catalyst.expressions.GeneratedClass$SpecificPredicate.subExpr_1$(Unknown Source)
                at org.apache.spark.sql.catalyst.expressions.GeneratedClass$SpecificPredicate.eval(Unknown Source)
                at org.apache.spark.sql.execution.FilterExec.$anonfun$doExecute$3(basicPhysicalOperators.scala:276)
                at org.apache.spark.sql.execution.FilterExec.$anonfun$doExecute$3$adapted(basicPhysicalOperators.scala:275)
                at scala.collection.Iterator$$anon$12.hasNext(Iterator.scala:515)
                at org.apache.spark.sql.catalyst.expressions.GeneratedClass$GeneratedIteratorForCodegenStage2.processNext(Unknown Source)
                at org.apache.spark.sql.execution.BufferedRowIterator.hasNext(BufferedRowIterator.java:43)
                at org.apache.spark.sql.execution.WholeStageCodegenExec$$anon$1.hasNext(WholeStageCodegenExec.scala:760)
                at scala.collection.Iterator$$anon$10.hasNext(Iterator.scala:460)
                at org.apache.spark.sql.catalyst.expressions.GeneratedClass$GeneratedIteratorForCodegenStage3.processNext(Unknown Source)
                at org.apache.spark.sql.execution.BufferedRowIterator.hasNext(BufferedRowIterator.java:43)
                at org.apache.spark.sql.execution.WholeStageCodegenExec$$anon$1.hasNext(WholeStageCodegenExec.scala:760)
                at org.apache.spark.sql.execution.SparkPlan.$anonfun$getByteArrayRdd$1(SparkPlan.scala:388)
                at org.apache.spark.rdd.RDD.$anonfun$mapPartitionsInternal$2(RDD.scala:888)
                at org.apache.spark.rdd.RDD.$anonfun$mapPartitionsInternal$2$adapted(RDD.scala:888)
                at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:52)
                at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:364)
                at org.apache.spark.rdd.RDD.iterator(RDD.scala:328)
                at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:92)
                at org.apache.spark.TaskContext.runTaskWithListeners(TaskContext.scala:161)
                at org.apache.spark.scheduler.Task.run(Task.scala:139)
                at org.apache.spark.executor.Executor$TaskRunner.$anonfun$run$3(Executor.scala:554)
                at org.apache.spark.util.Utils$.tryWithSafeFinally(Utils.scala:1529)
                at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:557)
                at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128)
                at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628)
                ... 1 more

测试数据与Schema

测试数据:

{
    'id': '2',
    'arr2': [1, 1, 1, 2],
    'arr1': [0.0, 1.0, 1.0, 1.0]
}

Schema:

root
|-- id: string (nullable = true)
|-- arr2: array (nullable = true)
|    |-- element: long (containsNull = true)
|-- arr1: array (nullable = true)
|    |-- element: double (containsNull = true)

执行代码

import pyspark.sql.functions as F
from pyspark.sql import SparkSession

spark = SparkSession.builder.master("local").appName("test-app").getOrCreate()
data = [
    {
        'id': '2',
        'arr2': [1, 1, 1, 2],
        'arr1': [0.0, 1.0, 1.0, 1.0]
    }
]
df = spark.createDataFrame(data)

df = df.withColumn('end_index', F.size('arr2') - F.lit(6))
df = df.filter(F.col('end_index') > 0) # 此处过滤end_index>0的行
df = df.withColumn("trimmed_arr2",F.slice(F.col('arr2'), start=F.lit(1), length=F.col('end_index')))
df = df.withColumn("avg_trimmed", F.expr('aggregate(trimmed_arr2, 0L, (acc,x) -> acc+x, acc -> acc / end_index)'))
df = df.filter(F.col('avg_trimmed') > 30)
df = df.withColumn('repeated_counts', F.size(F.array_distinct('trimmed_arr2')))
df = df.withColumn('ratio', F.col('repeated_counts') / F.size('trimmed_arr2'))
df = df.filter(F.col('ratio') > 0.6)
df.show(truncate=False, vertical=True)

异常出现的特殊现象:在执行df = df.filter(F.col('ratio') > 0.6)之前将DataFrame写入磁盘再读回,代码可正常执行无异常。

核心原因分析

对比排除PushDownPredicates规则前后的物理计划可以明确:

  • 未排除规则时:所有filter条件被合并后推到Scan操作之后,slice操作的length计算(size(arr2)-6)在end_index>0的filter之前就被执行。对于测试数据,size(arr2)=4,计算得end_index=4-6=-2,是负数,直接触发slice函数的参数校验异常。
  • 排除规则或写入磁盘后:end_index>0的filter会先执行,过滤掉不符合条件的行(测试数据中该行会被过滤),后续slice操作不会被执行,自然不会触发异常。

本质是Spark 3.2的Predicate PushDown优化存在逻辑漏洞,错误地将依赖于filter结果的slice操作下推到filter之前执行,导致原本应该被过滤的行提前进入slice计算,传入了非法的负数length。

排查方向与解决方案

排查方向

  • 验证优化规则的影响:通过添加配置spark.sql.optimizer.excludedRules=org.apache.spark.sql.catalyst.optimizer.PushDownPredicates禁用谓词下推,观察异常是否消失,确认问题根源。
  • 检查列依赖关系:确认计算slice length的列是否存在隐式依赖,是否优化器误判了计算顺序。
  • 版本兼容性测试:升级到PySpark 3.3及以上版本,验证该优化器bug是否已被修复。

解决方案

  1. 禁用特定优化规则:临时禁用PushDownPredicates,确保filter先于slice执行:
spark = SparkSession.builder.master("local").appName("test-app") \
    .config("spark.sql.optimizer.excludedRules", "org.apache.spark.sql.catalyst.optimizer.PushDownPredicates") \
    .getOrCreate()
  1. 强制物化打断优化链:在filter之后添加cache()或persist(),让Spark先执行filter并缓存结果,避免后续操作被下推:
df = df.filter(F.col('end_index') > 0).cache()
  1. 显式保证length非负:在slice操作中使用greatest函数强制length不小于0,即使优化器下推也不会触发参数校验异常:
df = df.withColumn("trimmed_arr2",F.slice(F.col('arr2'), start=F.lit(1), length=F.greatest(F.col('end_index'), F.lit(0))))

内容的提问来源于stack exchange,提问作者daniel bariudin

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.19 16:47:00