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

PySpark中Pandas UDF执行出现SparkException报错求助

解决PySpark Pandas UDF中的Arrow相关SparkException

咱们来搞定你遇到的这个问题!从错误栈来看,核心问题出在Apache Arrow的序列化/反序列化环节——java.nio.ByteBuffer.allocate抛出了非法参数异常,这大概率是Spark、Arrow、Pandas的版本不兼容,或者内存配置跟不上导致的。下面给你几个靠谱的解决方向:

1. 对齐Spark、Arrow、Pandas的版本(最关键)

PySpark的Pandas UDF靠Arrow在JVM和Python进程之间传数据,版本不对的话,两边的序列化格式不匹配,直接就会炸。给你一份常见的版本对应表:

  • Spark 2.4.x → Arrow 0.12.x~0.15.x,Pandas 0.23.x~0.25.x
  • Spark 3.0.x → Arrow 0.15.x~0.17.x,Pandas 1.0.x
  • Spark 3.1+ → 推荐Arrow 1.0+ + Pandas 1.1+

先检查你当前的版本:

pip show pyarrow pandas
spark-submit --version

如果版本不匹配,直接升级/降级到对应兼容版本:

pip install pyarrow==<对应版本号> pandas==<对应版本号>

2. 调整Arrow的内存相关配置

错误里的ByteBuffer分配失败,可能是Arrow要分配的内存超过了JVM或系统的限制。你可以在初始化SparkSession的时候加几个配置:

from pyspark.sql import SparkSession

spark = SparkSession.builder \
    .appName("PandasUDFTest") \
    .config("spark.sql.execution.arrow.maxRecordsPerBatch", "10000")  # 缩小每个Arrow批次的记录数,降低内存压力
    .config("spark.driver.memory", "8g")  # 根据你的机器情况调大Driver内存
    .config("spark.executor.memory", "8g")  # 同理调大Executor内存
    .getOrCreate()

减小批次记录数能直接降低单次需要分配的内存大小,大概率能绕过这个错误。

3. 临时禁用Arrow优化(应急方案)

如果上面的方法都不管用,你可以先关掉PySpark的Arrow优化,用传统序列化方式跑,虽然性能会差一点,但能先验证你的UDF逻辑对不对:

spark.conf.set("spark.sql.execution.arrow.enabled", "false")

之后再运行你的代码,如果不报错,那肯定就是Arrow相关的兼容性或内存问题了。

4. 检查数据里的特殊值

虽然你的Schema显示a列是double类型,但如果数据里有无穷大(inf)、NaN或者超大数值,Arrow序列化的时候也可能出问题。先检查一下:

# 查看a列的统计信息
df.select("a").describe().show()
# 统计特殊值的数量
df.filter(df.a.isin([float('inf'), float('-inf'), float('nan')])).count()

如果有这类值,先清洗一下(比如替换成合理值或者过滤掉),再跑UDF试试。


内容的提问来源于stack exchange,提问作者Prudhvi Raju Srivatsavaya

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.14 09:11:27