Spark3.4+Java17升级后无法读取AWS S3中Parquet文件求助
问题描述
我们正将Spark作业从Java 11 + Spark 3.0升级至Java 17 + Spark 3.4,此前在旧环境下运行正常,但升级后应用无法读取AWS S3中的Parquet文件,容器退出码为50,错误堆栈如下:
container name: app-executor container image: imagepath container state: terminated container started at: 2023-12-09T14:30:43Z container finished at: 2023-12-09T14:30:47Z exit code: 50 termination reason: Error Driver stacktrace: at org.apache.spark.scheduler.DAGScheduler.failJobAndIndependentStages(DAGScheduler.scala:2785) at org.apache.spark.scheduler.DAGScheduler.$anonfun$abortStage$2(DAGScheduler.scala:2721) at org.apache.spark.scheduler.DAGScheduler.$anonfun$abortStage$2$adapted(DAGScheduler.scala:2720) at scala.collection.mutable.ResizableArray.foreach(ResizableArray.scala:62) at scala.collection.mutable.ResizableArray.foreach$(ResizableArray.scala:55) at scala.collection.mutable.ArrayBuffer.foreach(ArrayBuffer.scala:49) at org.apache.spark.scheduler.DAGScheduler.abortStage(DAGScheduler.scala:2720) at org.apache.spark.scheduler.DAGScheduler.$anonfun$handleTaskSetFailed$1(DAGScheduler.scala:1206) at org.apache.spark.scheduler.DAGScheduler.$anonfun$handleTaskSetFailed$1$adapted(DAGScheduler.scala:1206) at scala.Option.foreach(Option.scala:407) at org.apache.spark.scheduler.DAGScheduler.handleTaskSetFailed(DAGScheduler.scala:1206) at org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.doOnReceive(DAGScheduler.scala:2984) at org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.onReceive(DAGScheduler.scala:2923) at org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.onReceive(DAGScheduler.scala:2912) at org.apache.spark.util.EventLoop$$anon$1.run(EventLoop.scala:49) at org.apache.spark.scheduler.DAGScheduler.runJob(DAGScheduler.scala:971) at org.apache.spark.SparkContext.runJob(SparkContext.scala:2263) at org.apache.spark.SparkContext.runJob(SparkContext.scala:2284) at org.apache.spark.SparkContext.runJob(SparkContext.scala:2303) at org.apache.spark.SparkContext.runJob(SparkContext.scala:2328) at org.apache.spark.rdd.RDD.$anonfun$collect$1(RDD.scala:1019) at org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:151) at org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:112) at org.apache.spark.rdd.RDD.withScope(RDD.scala:405) at org.apache.spark.rdd.RDD.collect(RDD.scala:1018) at org.apache.spark.sql.execution.datasources.SchemaMergeUtils$.mergeSchemasInParallel(SchemaMergeUtils.scala:73) at org.apache.spark.sql.execution.datasources.parquet.ParquetFileFormat$.mergeSchemasInParallel(ParquetFileFormat.scala:476) at org.apache.spark.sql.execution.datasources.parquet.ParquetUtils$.inferSchema(ParquetUtils.scala:132) at org.apache.spark.sql.execution.datasources.parquet.ParquetFileFormat.inferSchema(ParquetFileFormat.scala:78) at org.apache.spark.sql.execution.datasources.DataSource.$anonfun$getOrInferFileFormatSchema$11(DataSource.scala:208) at scala.Option.orElse(Option.scala:447) at org.apache.spark.sql.execution.datasources.DataSource.getOrInferFileFormatSchema(DataSource.scala:205) at org.apache.spark.sql.execution.datasources.DataSource.resolveRelation(DataSource.scala:407) at org.apache.spark.sql.DataFrameReader.loadV1Source(DataFrameReader.scala:229) at org.apache.spark.sql.DataFrameReader.$anonfun$load$2(DataFrameReader.scala:211) at scala.Option.getOrElse(Option.scala:189) at org.apache.spark.sql.DataFrameReader.load(DataFrameReader.scala:211) at org.apache.spark.sql.DataFrameReader.parquet(DataFrameReader.scala:563) at org.apache.spark.sql.DataFrameReader.parquet(DataFrameReader.scala:562) at app.spark.reader.ParquetReader.read(ParquetReader.java:58) at app.spark.reader.SparkJob.run(SparkJob.java:49) at app.spark.reader.SparkJobRunner.run(SparkJobRunner.java:43) at org.springframework.boot.SpringApplication.callRunner(SpringApplication.java:768) at org.springframework.boot.SpringApplication.callRunners(SpringApplication.java:752) at org.springframework.boot.SpringApplication.run(SpringApplication.java:314) at org.springframework.boot.SpringApplication.run(SpringApplication.java:1303) at org.springframework.boot.SpringApplication.run(SpringApplication.java:1292) at net.jpmchase.gbpde.pir.eod.Application.main(Application.java:20) at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke0(Native Method) at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:77) at java.base/jdk.internal.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43) at java.base/java.lang.reflect.Method.invoke(Method.java:568) at org.apache.spark.deploy.JavaMainApplication.start(SparkApplication.scala:52) at org.apache.spark.deploy.SparkSubmit.org$apache$spark$deploy$SparkSubmit$$runMain(SparkSubmit.scala:1020) at org.apache.spark.deploy.SparkSubmit.doRunMain$1(SparkSubmit.scala:192) at org.apache.spark.deploy.SparkSubmit.submit(SparkSubmit.scala:215) at org.apache.spark.deploy.SparkSubmit.doSubmit(SparkSubmit.scala:91) at org.apache.spark.deploy.SparkSubmit$$anon$2.doSubmit(SparkSubmit.scala:1111) at org.apache.spark.deploy.SparkSubmit$.main(SparkSubmit.scala:1120) at org.apache.spark.deploy.SparkSubmit.main(SparkSubmit.scala)
排查方向与解决方案
从堆栈看,错误发生在Parquet schema并行合并阶段,结合Java17和Spark3.4的升级背景,重点排查以下几点:
1. S3客户端依赖冲突或版本不兼容
Spark 3.4默认使用的AWS SDK版本与Spark 3.0不同,Java17对依赖兼容性要求更严格:
- 检查项目中是否引入旧版本AWS SDK(如
aws-java-sdk-s3),建议统一使用Spark 3.4内置的AWS SDK版本,排除冲突依赖。 - 确保配置正确的S3文件系统实现,Spark 3.4推荐使用
org.apache.hadoop.fs.s3a.S3AFileSystem,需设置参数:spark.hadoop.fs.s3a.impl=org.apache.hadoop.fs.s3a.S3AFileSystem
2. Java 17模块系统权限问题
Java17的模块化可能限制Spark或AWS SDK的反射访问,导致Parquet解析出错:
- 添加JVM参数开放相关模块权限:
--add-opens java.base/java.lang=ALL-UNNAMED --add-opens java.base/java.nio=ALL-UNNAMED --add-opens java.base/java.util=ALL-UNNAMED - 检查是否有Java安全策略文件限制S3访问权限,必要时调整策略。
3. Parquet schema合并兼容性问题
Spark 3.4对Parquet schema合并逻辑做了优化,可能与旧版本生成的Parquet文件存在兼容性问题:
- 尝试禁用自动schema合并,指定明确的schema读取文件:
StructType schema = ...; // 提前定义好的schema spark.read().schema(schema).parquet("s3a://path/to/files"); - 检查S3中的Parquet文件是否存在schema不一致情况,Spark 3.4对schema冲突的处理更严格。
4. Spark executor资源不足
容器快速退出(启动4秒后结束)可能是executor内存不足导致OOM:
- 增加executor的内存和CPU配置,例如:
spark.executor.memory=4g spark.executor.cores=2 - 查看容器日志中是否有OOM相关的隐藏错误信息。
5. AWS权限配置问题
升级后可能因环境变量或IAM角色变化导致S3访问权限不足:
- 验证运行作业的IAM角色是否有S3对象的
GetObject权限。 - 检查Spark配置中的AWS凭证是否正确,避免硬编码凭证,推荐使用IAM角色或环境变量。
内容的提问来源于stack exchange,提问作者Raju V
相关产品推荐
相关产品推荐

