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

使用Pandas分组Map UDF聚合DataFrame时触发Java错误

环境信息

  • AWS EMR版本标签:6.1.0
  • Spark版本:3.0.0
  • Pandas版本:1.1.0
  • Pyarrow版本:0.15.1(曾升级至1.0.1、12.0.0测试)
  • Python版本:3.7.16

问题描述

在集群关联的Jupyter Notebook中使用分组Map类型的Pandas UDF处理DataFrame时,出现以下异常:

  1. 基于my_schema定义分组Map函数my_function:
@pandas_udf(my_schema, PandasUDFType.GROUPED_MAP)
def my_function(my_df):
    # 业务逻辑处理
    return my_df
  1. 调用代码:
my_df_new = (my_df.drop('some_column').groupby('some_other_column').apply(my_function))
  1. 现象:将my_df_new创建临时视图后,select * from my_df_new可正常返回结果,但执行聚合查询(如select count(*) from my_df_new)时触发Java错误。

已尝试的无效解决方案

  • 修改Spark Session配置:
    • spark.driver.maxResultSize: "0"
    • spark.sql.execution.arrow.pyspark.enabled: "true"
    • spark.sql.execution.pandas.udf.buffer.size: "2000000000"
    • spark.sql.execution.arrow.maxRecordsPerBatch: "33554432"
  • 升级Pyarrow版本至1.0.1和12.0.0

错误信息

调用o147.showString时发生错误。
: org.apache.spark.SparkException: 因阶段失败导致作业中止:阶段20.0中的任务151失败4次,最近一次失败:丢失阶段20.0中的任务151.3(TID 14659,ip-xx-xxx-xx-xxx.my_domain.com,执行器47):java.lang.IndexOutOfBoundsException: 索引:0,长度:1073741824(预期范围:(0, 0))
    at io.netty.buffer.ArrowBuf.checkIndex(ArrowBuf.java:716)
    at io.netty.buffer.ArrowBuf.setBytes(ArrowBuf.java:954)
    at org.apache.arrow.vector.BaseVariableWidthVector.reallocDataBuffer(BaseVariableWidthVector.java:508)
    at org.apache.arrow.vector.BaseVariableWidthVector.handleSafe(BaseVariableWidthVector.java:1239)
    at org.apache.arrow.vector.BaseVariableWidthVector.setSafe(BaseVariableWidthVector.java:1066)
    at org.apache.spark.sql.execution.arrow.StringWriter.setValue(ArrowWriter.scala:248)
    at org.apache.spark.sql.execution.arrow.ArrowFieldWriter.write(ArrowWriter.scala:127)
    at org.apache.spark.sql.execution.arrow.ArrayWriter.setValue(ArrowWriter.scala:300)
    at org.apache.spark.sql.execution.arrow.ArrowFieldWriter.write(ArrowWriter.scala:127)
    at org.apache.spark.sql.execution.arrow.ArrowWriter.write(ArrowWriter.scala:92)
    at org.apache.spark.sql.execution.python.ArrowPythonRunner$$anon$1.$anonfun$writeIteratorToStream$1(ArrowPythonRunner.scala:90)
    at scala.runtime.java8.JFunction0$mcV$sp.apply(JFunction0$mcV$sp.java:23)
    at org.apache.spark.util.Utils$.tryWithSafeFinally(Utils.scala:1377)
    at org.apache.spark.sql.execution.python.ArrowPythonRunner$$anon$1.writeIteratorToStream(ArrowPythonRunner.scala:101)
    at org.apache.spark.api.python.BasePythonRunner$WriterThread.$anonfun$run$1(PythonRunner.scala:383)
    at org.apache.spark.util.Utils$.logUncaughtExceptions(Utils.scala:1932)
    at org.apache.spark.api.python.BasePythonRunner$WriterThread.run(PythonRunner.scala:218)

可行解决思路

  1. 校验UDF返回结果与schema的一致性

    • 确认my_function返回的DataFrame列名、数据类型完全匹配my_schema,聚合查询会触发更严格的schema校验,select *可能未暴露隐性类型不匹配问题。
    • 检查是否存在返回空DataFrame的分组,空DataFrame的schema可能与定义的my_schema存在差异,可在UDF中对空分组做兜底处理(返回符合schema的空DataFrame)。
  2. 禁用Arrow优化

    • 错误源于Arrow数据传输环节,尝试设置spark.sql.execution.arrow.pyspark.enabled=false,强制Spark用传统序列化方式处理Pandas UDF的数据交互,避开Arrow的内存缓冲问题。
  3. 调整分组粒度与数据分区

    • 若some_other_column的分组基数过大或部分分组数据量异常庞大,会导致Arrow处理单批数据时内存溢出:
      • 增加分组键拆分大分组;
      • 调整spark.sql.shuffle.partitions增加shuffle分区数,减少每个分区处理的数据量;
      • 对数据抽样,验证是否是特定大分组导致的问题。
  4. 排查UDF业务逻辑

    • 检查my_function是否生成了异常长度的字符串或数组(如1GB级别的字符串),这类数据会触发Arrow内存索引越界;
    • 清理或转换UDF中的NaN、None等特殊值,这些值在Arrow序列化时可能引发异常。
  5. 回退Pyarrow至兼容版本

    • Spark 3.0.0官方推荐Pyarrow 0.15.x版本,过高版本可能存在兼容性问题,回退到0.15.1并确保集群所有节点(driver和executor)的Pyarrow版本一致。
  6. 调整Spark内存配置

    • 增加spark.executor.memory和spark.executor.memoryOverhead,避免executor处理大分组时内存不足导致Arrow缓冲异常。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.21 11:30:44