EMR中PySpark写入Parquet到S3时出现InvalidPart错误
问题分析与解决方案
错误核心原因
偶发的InvalidPart错误是Spark向S3写入Parquet时,多部分上传的分片丢失或ETag不匹配导致的。权限已确认正常,问题大概率源于网络波动、S3最终一致性延迟,或是Spark与S3交互的配置优化不足。
具体解决方案
使用优化的S3文件提交器
切换到EMR推荐的提交策略,避免多部分上传的一致性问题:# 启用DIRECT提交策略(EMR 5.20+支持) spark.conf.set("spark.sql.s3.commitStrategy", "DIRECT") # 或使用S3A目录提交器 spark.conf.set("spark.hadoop.fs.s3a.committer.name", "directory") spark.conf.set("spark.hadoop.fs.s3a.committer.staging.conflict-mode", "replace")同时开启EMRFS强一致性:
spark.conf.set("emrfs.s3.consistent", "true")调整多部分上传参数
增大分片大小,减少分片数量,降低网络波动带来的分片丢失概率:# 设置分片大小为128MB(默认是10MB) spark.conf.set("spark.hadoop.fs.s3a.multipart.size", "134217728") # 设置触发多部分上传的阈值为128MB spark.conf.set("spark.hadoop.fs.s3a.multipart.threshold", "134217728")增加重试次数
针对偶发的网络错误,提高任务和S3操作的重试次数:# 提高任务最大失败重试次数 spark.conf.set("spark.task.maxFailures", "4") # 增加S3操作的重试次数 spark.conf.set("spark.sql.s3.retry.maxAttempts", "10") spark.conf.set("spark.hadoop.fs.s3a.retry.max", "10")升级EMR集群版本
旧版本EMR的S3客户端可能存在多部分上传的bug,建议升级到EMR 6.x及以上版本,获得更稳定的S3交互支持。避免并发写入冲突
确保不同任务写入独立的S3路径,比如按时间戳生成临时目录,完成后再合并到目标路径,避免同一路径下的并发写入导致分片冲突。
原始错误信息
Py4JJavaError: An error occurred while calling o426.parquet. : org.apache.spark.SparkException: Job aborted due to stage failure: Authorized committer (attemptNumber=0, stage=17, partition=264) failed; but task commit success, data duplication may happen. reason=ExceptionFailure(org.apache.spark.SparkException,[TASK_WRITE_FAILED] Task failed while writing rows to s3://bucket/path.,[Ljava.lang.StackTraceElement;@397e21a9,org.apache.spark.SparkException: [TASK_WRITE_FAILED] Task failed while writing rows to s3://bucket/path. at org.apache.spark.sql.errors.QueryExecutionErrors$.taskFailedWhileWritingRowsError(QueryExecutionErrors.scala:789) at org.apache.spark.sql.execution.datasources.FileFormatWriter$.executeTask(FileFormatWriter.scala:421) at org.apache.spark.sql.execution.datasources.WriteFilesExec.$anonfun$doExecuteWrite$1(WriteFiles.scala:100) at org.apache.spark.rdd.RDD.$anonfun$mapPartitionsInternal$2(RDD.scala:888) at org.apache.spark.rdd.RDD.$anonfun$mapPartitionsInternal$2$adapted(RDD.scala:888) at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:52) at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:364) at org.apache.spark.rdd.RDD.iterator(RDD.scala:328) at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:92) at org.apache.spark.TaskContext.runTaskWithListeners(TaskContext.scala:161) at org.apache.spark.scheduler.Task.run(Task.scala:141) at org.apache.spark.executor.Executor$TaskRunner.$anonfun$run$3(Executor.scala:563) at org.apache.spark.util.Utils$.tryWithSafeFinally(Utils.scala:1541) at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:566) at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149) at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624) at java.lang.Thread.run(Thread.java:750) Caused by: com.amazon.ws.emr.hadoop.fs.shaded.com.amazonaws.services.s3.model.AmazonS3Exception: One or more of the specified parts could not be found. The part may not have been uploaded, or the specified entity tag may not match the part's entity tag. (Service: Amazon S3; Status Code: 400; Error Code: InvalidPart; Request ID: 0RGE13WMZ76BMPW6; S3 Extended Request ID: up90NKdAy7UIp3Rep2+J293TUhfFcno8iG/Y7Qr8uZOLMMzrQAwrZrfKojzKsq5iKiuGPQLz9/g=; Proxy: null), S3 Extended Request ID: up90NKdAy7UIp3Rep2+J293TUhfFcno8iG/Y7Qr8uZOLMMzrQAwrZrfKojzKsq5iKiuGPQLz9/g=
内容的提问来源于stack exchange,提问作者Shadowtrooper
相关产品推荐
相关产品推荐

