Spark写入S3时触发java.lang.ArrayIndexOutOfBoundsException:26错误求助
问题描述
使用Spark-Snowflake连接器从Snowflake表读取数据,执行转换后写入S3时,执行代码df.write.partitionBy(partitionColumn).mode(SaveMode.Append).parquet(s3Destination)触发java.lang.ArrayIndexOutOfBoundsException: 26错误,完整堆栈信息如下:
java.lang.ArrayIndexOutOfBoundsException: 26 at org.apache.spark.sql.types.StructType.apply(StructType.scala:414) at org.apache.spark.sql.catalyst.optimizer.ObjectSerializerPruning$$anonfun$apply$4.$anonfun$applyOrElse$3(objects.scala:224) at scala.collection.TraversableLike.$anonfun$map$1(TraversableLike.scala:238) at scala.collection.immutable.List.foreach(List.scala:392) at scala.collection.TraversableLike.map(TraversableLike.scala:238) at scala.collection.TraversableLike.map$(TraversableLike.scala:231) at scala.collection.immutable.List.map(List.scala:298) at org.apache.spark.sql.catalyst.optimizer.ObjectSerializerPruning$$anonfun$apply$4.applyOrElse(objects.scala:223) at org.apache.spark.sql.catalyst.optimizer.ObjectSerializerPruning$$anonfun$apply$4.applyOrElse(objects.scala:211) at org.apache.spark.sql.catalyst.trees.TreeNode.$anonfun$transformDown$1(TreeNode.scala:328) at org.apache.spark.sql.catalyst.trees.CurrentOrigin$.withOrigin(TreeNode.scala:74) at org.apache.spark.sql.catalyst.trees.TreeNode.transformDown(TreeNode.scala:328) at org.apache.spark.sql.catalyst.plans.logical.LogicalPlan.org$apache$spark$sql$catalyst$plans$logical$AnalysisHelper$$super$transformDown(LogicalPlan.scala:29) at org.apache.spark.sql.catalyst.plans.logical.AnalysisHelper.transformDown(AnalysisHelper.scala:171) at org.apache.spark.sql.catalyst.plans.logical.AnalysisHelper.transformDown$(AnalysisHelper.scala:169) at org.apache.spark.sql.catalyst.plans.logical.LogicalPlan.transformDown(LogicalPlan.scala:29) at org.apache.spark.sql.catalyst.plans.logical.LogicalPlan.transformDown(LogicalPlan.scala:29) at org.apache.spark.sql.catalyst.trees.TreeNode.$anonfun$transformDown$3(TreeNode.scala:333) at org.apache.spark.sql.catalyst.trees.TreeNode.applyFunctionIfChanged$1(TreeNode.scala:387) at org.apache.spark.sql.catalyst.trees.TreeNode.$anonfun$mapChildren$1(TreeNode.scala:423) at org.apache.spark.sql.catalyst.trees.TreeNode.mapProductIterator(TreeNode.scala:255) at org.apache.spark.sql.catalyst.trees.TreeNode.mapChildren(TreeNode.scala:421) at org.apache.spark.sql.catalyst.trees.TreeNode.mapChildren(TreeNode.scala:369) at org.apache.spark.sql.catalyst.trees.TreeNode.transformDown(TreeNode.scala:333) at org.apache.spark.sql.catalyst.plans.logical.LogicalPlan.org$apache$spark$sql$catalyst$plans$logical$AnalysisHelper$$super$transformDown(LogicalPlan.scala:29) at org.apache.spark.sql.catalyst.plans.logical.AnalysisHelper.transformDown(AnalysisHelper.scala:171) at org.apache.spark.sql.catalyst.plans.logical.AnalysisHelper.transformDown$(AnalysisHelper.scala:169) at org.apache.spark.sql.catalyst.plans.logical.LogicalPlan.transformDown(LogicalPlan.scala:29) at org.apache.spark.sql.catalyst.plans.logical.LogicalPlan.transformDown(LogicalPlan.scala:29) at org.apache.spark.sql.catalyst.trees.TreeNode.$anonfun$transformDown$3(TreeNode.scala:333) at org.apache.spark.sql.catalyst.trees.TreeNode.applyFunctionIfChanged$1(TreeNode.scala:387) at org.apache.spark.sql.catalyst.trees.TreeNode.$anonfun$mapChildren$1(TreeNode.scala:423) at org.apache.spark.sql.catalyst.trees.TreeNode.mapProductIterator(TreeNode.scala:255) at org.apache.spark.sql.catalyst.trees.TreeNode.mapChildren(TreeNode.scala:421) at org.apache.spark.sql.catalyst.trees.TreeNode.mapChildren(TreeNode.scala:369) at org.apache.spark.sql.catalyst.trees.TreeNode.transformDown(TreeNode.scala:333) at org.apache.spark.sql.catalyst.plans.logical.LogicalPlan.org$apache$spark$sql$catalyst$plans$logical$AnalysisHelper$$super$transformDown(LogicalPlan.scala:29) at org.apache.spark.sql.catalyst.plans.logical.AnalysisHelper.transformDown(AnalysisHelper.scala:171) at org.apache.spark.sql.catalyst.plans.logical.AnalysisHelper.transformDown$(AnalysisHelper.scala:169) at org.apache.spark.sql.catalyst.plans.logical.LogicalPlan.transformDown(LogicalPlan.scala:29) at org.apache.spark.sql.catalyst.plans.logical.LogicalPlan.transformDown(LogicalPlan.scala:29) at org.apache.spark.sql.catalyst.trees.TreeNode.transform(TreeNode.scala:317) at org.apache.spark.sql.catalyst.optimizer.ObjectSerializerPruning$.apply(objects.scala:211) at org.apache.spark.sql.catalyst.optimizer.ObjectSerializerPruning$.apply(objects.scala:121) at org.apache.spark.sql.catalyst.rules.RuleExecutor.$anonfun$execute$2(RuleExecutor.scala:216) at scala.collection.IndexedSeqOptimized.foldLeft(IndexedSeqOptimized.scala:60) at scala.collection.IndexedSeqOptimized.foldLeft$(IndexedSeqOptimized.scala:68) at scala.collection.mutable.WrappedArray.foldLeft(WrappedArray.scala:38) at org.apache.spark.sql.catalyst.rules.RuleExecutor.$anonfun$execute$1(RuleExecutor.scala:213) at org.apache.spark.sql.catalyst.rules.RuleExecutor.$anonfun$execute$1$adapted(RuleExecutor.scala:205) at scala.collection.immutable.List.foreach(List.scala:392) at org.apache.spark.sql.catalyst.rules.RuleExecutor.execute(RuleExecutor.scala:205) at org.apache.spark.sql.catalyst.rules.RuleExecutor.$anonfun$executeAndTrack$1(RuleExecutor.scala:183) at org.apache.spark.sql.catalyst.QueryPlanningTracker$.withTracker(QueryPlanningTracker.scala:107) at org.apache.spark.sql.catalyst.rules.RuleExecutor.executeAndTrack(RuleExecutor.scala:183) at org.apache.spark.sql.execution.QueryExecution.$anonfun$optimizedPlan$1(QueryExecution.scala:87) at org.apache.spark.sql.catalyst.QueryPlanningTracker.measurePhase(QueryPlanningTracker.scala:192) at org.apache.spark.sql.execution.QueryExecution.$anonfun$executePhase$1(QueryExecution.scala:163) at org.apache.spark.sql.SparkSession.withActive(SparkSession.scala:772) at org.apache.spark.sql.execution.QueryExecution.executePhase(QueryExecution.scala:163) at org.apache.spark.sql.execution.QueryExecution.optimizedPlan$lzycompute(QueryExecution.scala:84) at org.apache.spark.sql.execution.QueryExecution.optimizedPlan(QueryExecution.scala:84) at org.apache.spark.sql.execution.QueryExecution.$anonfun$writePlans$4(QueryExecution.scala:243) at org.apache.spark.sql.catalyst.plans.QueryPlan$.append(QueryPlan.scala:487) at org.apache.spark.sql.execution.QueryExecution.writePlans(QueryExecution.scala:243) at org.apache.spark.sql.execution.QueryExecution.toString(QueryExecution.scala:260) at org.apache.spark.sql.execution.QueryExecution.org$apache$spark$sql$execution$QueryExecution$$explainString(QueryExecution.scala:216) at org.apache.spark.sql.execution.QueryExecution.explainString(QueryExecution.scala:195) at org.apache.spark.sql.execution.SQLExecution$.executeQuery$1(SQLExecution.scala:102) at org.apache.spark.sql.execution.SQLExecution$.$anonfun$withNewExecutionId$6(SQLExecution.scala:135) at org.apache.spark.sql.catalyst.QueryPlanningTracker$.withTracker(QueryPlanningTracker.scala:107) at org.apache.spark.sql.execution.SQLExecution$.withTracker(SQLExecution.scala:232) at org.apache.spark.sql.execution.SQLExecution$.$anonfun$withNewExecutionId$5(SQLExecution.scala:135) at org.apache.spark.sql.execution.SQLExecution$.withSQLConfPropagated(SQLExecution.scala:253) at org.apache.spark.sql.execution.SQLExecution$.$anonfun$withNewExecutionId$1(SQLExecution.scala:134) at org.apache.spark.sql.SparkSession.withActive(SparkSession.scala:772) at org.apache.spark.sql.execution.SQLExecution$.withNewExecutionId(SQLExecution.scala:68) at org.apache.spark.sql.DataFrameWriter.runCommand(DataFrameWriter.scala:989) at org.apache.spark.sql.DataFrameWriter.saveToV1Source(DataFrameWriter.scala:438) at org.apache.spark.sql.DataFrameWriter.saveInternal(DataFrameWriter.scala:415) at org.apache.spark.sql.DataFrameWriter.save(DataFrameWriter.scala:293) at org.apache.spark.sql.DataFrameWriter.parquet(DataFrameWriter.scala:874) at spark.ExportEvents$.writeToS3(ExportEvents.scala:315) at spark.ExportEvents$.$anonfun$processIncrementalLoad$1(ExportEvents.scala:257) at scala.collection.immutable.List.map(List.scala:286) at spark.ExportEvents$.processIncrementalLoad(ExportEvents.scala:185) at spark.ExportEvents$.main(ExportEvents.scala:75) at spark.ExportEvents.main(ExportEvents.scala) at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method) at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62) at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43) at java.lang.reflect.Method.invoke(Method.java:498) at org.apache.spark.deploy.JavaMainApplication.start(SparkApplication.scala:52) at org.apache.spark.deploy.SparkSubmit.org$apache$spark$deploy$SparkSubmit$$runMain(SparkSubmit.scala:959) at org.apache.spark.deploy.SparkSubmit.doRunMain$1(SparkSubmit.scala:180) at org.apache.spark.deploy.SparkSubmit.submit(SparkSubmit.scala:203) at org.apache.spark.deploy.SparkSubmit.doSubmit(SparkSubmit.scala:90) at org.apache.spark.deploy.SparkSubmit$$anon$2.doSubmit(SparkSubmit.scala:1038) at org.apache.spark.deploy.SparkSubmit$.main(SparkSubmit.scala:1047) at org.apache.spark.deploy.SparkSubmit.main(SparkSubmit.scala)
转换过程中修改了一个列名为小写,其余列均为大写,使用版本为Spark 3.1.1和spark-snowflake_2.12。
解决方案
1. 统一列名大小写并排查冲突
- 检查转换后的DataFrame,确保没有大小写重复的列(比如
COL_A和col_a同时存在),这类情况会导致优化器元数据混乱 - 显式设置列名大小写敏感配置,初始化SparkSession时添加:
或者将所有列名统一为大写/小写,避免混合模式。spark.conf.set("spark.sql.caseSensitive", "true")
2. 禁用问题优化规则
堆栈显示错误发生在ObjectSerializerPruning优化阶段,Spark 3.1.x存在该规则的已知bug,可临时禁用:
spark.conf.set("spark.sql.optimizer.excludedRules", "org.apache.spark.sql.catalyst.optimizer.ObjectSerializerPruning")
3. 校验DataFrame Schema一致性
- 读取Snowflake数据后执行
df.printSchema(),转换后再次打印schema,确认列数、列名、数据类型无异常 - 重点检查修改列名后的schema是否正确,避免出现索引越界的列定义
4. 升级Spark版本
Spark 3.1.1存在多个优化器相关bug,升级到3.1.3或3.2.x稳定版本可修复此类问题,同时确保spark-snowflake连接器版本与Spark版本兼容。
5. 校验Append模式下的Schema兼容性
如果是Append模式,检查S3中已有Parquet文件的schema与当前DataFrame的schema是否完全一致(包括列名大小写、数据类型、列顺序),schema不匹配会触发优化器处理异常。
内容的提问来源于stack exchange,提问作者Hemanth Annavarapu
相关产品推荐
相关产品推荐

