使用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时,出现以下异常:
- 基于
my_schema定义分组Map函数my_function:
@pandas_udf(my_schema, PandasUDFType.GROUPED_MAP) def my_function(my_df): # 业务逻辑处理 return my_df
- 调用代码:
my_df_new = (my_df.drop('some_column').groupby('some_other_column').apply(my_function))
- 现象:将
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)
可行解决思路
校验UDF返回结果与schema的一致性
- 确认
my_function返回的DataFrame列名、数据类型完全匹配my_schema,聚合查询会触发更严格的schema校验,select *可能未暴露隐性类型不匹配问题。 - 检查是否存在返回空DataFrame的分组,空DataFrame的schema可能与定义的
my_schema存在差异,可在UDF中对空分组做兜底处理(返回符合schema的空DataFrame)。
- 确认
禁用Arrow优化
- 错误源于Arrow数据传输环节,尝试设置
spark.sql.execution.arrow.pyspark.enabled=false,强制Spark用传统序列化方式处理Pandas UDF的数据交互,避开Arrow的内存缓冲问题。
- 错误源于Arrow数据传输环节,尝试设置
调整分组粒度与数据分区
- 若
some_other_column的分组基数过大或部分分组数据量异常庞大,会导致Arrow处理单批数据时内存溢出:- 增加分组键拆分大分组;
- 调整
spark.sql.shuffle.partitions增加shuffle分区数,减少每个分区处理的数据量; - 对数据抽样,验证是否是特定大分组导致的问题。
- 若
排查UDF业务逻辑
- 检查
my_function是否生成了异常长度的字符串或数组(如1GB级别的字符串),这类数据会触发Arrow内存索引越界; - 清理或转换UDF中的NaN、None等特殊值,这些值在Arrow序列化时可能引发异常。
- 检查
回退Pyarrow至兼容版本
- Spark 3.0.0官方推荐Pyarrow 0.15.x版本,过高版本可能存在兼容性问题,回退到0.15.1并确保集群所有节点(driver和executor)的Pyarrow版本一致。
调整Spark内存配置
- 增加
spark.executor.memory和spark.executor.memoryOverhead,避免executor处理大分组时内存不足导致Arrow缓冲异常。
- 增加
内容的提问来源于stack exchange,提问作者fjjones88
相关产品推荐
相关产品推荐

