使用PySpark读取Minio数据时触发NoClassDefFoundError求助
PySpark读取Minio时出现NoClassDefFoundError:org/apache/hadoop/fs/impl/prefetch/PrefetchingStatistics
环境
- 设备:Macbook Air M1
- Minio:本地9000端口运行
- Spark:使用
apache/spark-py官方镜像
复现步骤
- 本地启动Minio实例(端口9000)
- 使用以下Dockerfile构建镜像:
FROM apache/spark-py COPY ./hadoop-aws-3.3.6.jar /opt/spark/jars/hadoop-aws-3.3.6.jar COPY ./aws-java-sdk-bundle-1.12.367.jar /opt/spark/jars/aws-java-sdk-bundle-1.12.367.jar COPY . .
- 运行容器:
docker run --network="host" -it <image-name> /opt/spark/bin/pyspark
- 在PySpark中执行以下代码:
from pyspark.sql import SparkSession spark = SparkSession.builder.getOrCreate() spark.sparkContext._jsc.hadoopConfiguration().set("fs.s3a.access.key", "minio") spark.sparkContext._jsc.hadoopConfiguration().set("fs.s3a.secret.key", "minio123") spark.sparkContext._jsc.hadoopConfiguration().set("fs.s3a.path.style.access", "true") spark.sparkContext._jsc.hadoopConfiguration().set("fs.s3a.impl", "org.apache.hadoop.fs.s3a.S3AFileSystem") spark.sparkContext._jsc.hadoopConfiguration().set("fs.s3a.endpoint", "http://localhost:9000") spark.sparkContext._jsc.hadoopConfiguration().set("fs.s3a.connection.ssl.enabled", "false") df = spark.read.json('s3a://lake-test/data.json')
错误信息
Traceback (most recent call last): File "<stdin>", line 1, in <module> File "/opt/spark/python/pyspark/sql/readwriter.py", line 418, in json return self._df(self._jreader.json(self._spark._sc._jvm.PythonUtils.toSeq(path))) File "/opt/spark/python/lib/py4j-0.10.9.7-src.zip/py4j/java_gateway.py", line 1322, in __call__ File "/opt/spark/python/pyspark/errors/exceptions/captured.py", line 169, in deco return f(*a, **kw) File "/opt/spark/python/lib/py4j-0.10.9.7-src.zip/py4j/protocol.py", line 326, in get_return_value py4j.protocol.Py4JJavaError: An error occurred while calling o45.json. : java.lang.NoClassDefFoundError: org/apache/hadoop/fs/impl/prefetch/PrefetchingStatistics at java.base/java.lang.ClassLoader.defineClass1(Native Method) at java.base/java.lang.ClassLoader.defineClass(Unknown Source) at java.base/java.security.SecureClassLoader.defineClass(Unknown Source) at java.base/jdk.internal.loader.BuiltinClassLoader.defineClass(Unknown Source) at java.base/jdk.internal.loader.BuiltinClassLoader.findClassOnClassPathOrNull(Unknown Source) at java.base/jdk.internal.loader.BuiltinClassLoader.loadClassOrNull(Unknown Source) at java.base/jdk.internal.loader.BuiltinClassLoader.loadClass(Unknown Source) at java.base/jdk.internal.loader.ClassLoaders$AppClassLoader.loadClass(Unknown Source) at java.base/java.lang.ClassLoader.loadClass(Unknown Source) at org.apache.hadoop.fs.s3a.S3AFileSystem.initialize(S3AFileSystem.java:519) at org.apache.hadoop.fs.FileSystem.createFileSystem(FileSystem.java:3469) at org.apache.hadoop.fs.FileSystem.access$300(FileSystem.java:174) at org.apache.hadoop.fs.FileSystem$Cache.getInternal(FileSystem.java:3574) at org.apache.hadoop.fs.FileSystem$Cache.get(FileSystem.java:3521) at org.apache.hadoop.fs.FileSystem.get(FileSystem.java:540) at org.apache.hadoop.fs.Path.getFileSystem(Path.java:365) at org.apache.spark.sql.execution.streaming.FileStreamSink$.hasMetadata(FileStreamSink.scala:53) at org.apache.spark.sql.execution.datasources.DataSource.resolveRelation(DataSource.scala:366) 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.json(DataFrameReader.scala:362) at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke0(Native Method) at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke(Unknown Source) at java.base/jdk.internal.reflect.DelegatingMethodAccessorImpl.invoke(Unknown Source) at java.base/java.lang.reflect.Method.invoke(Unknown Source) at py4j.reflection.MethodInvoker.invoke(MethodInvoker.java:244) at py4j.reflection.ReflectionEngine.invoke(ReflectionEngine.java:374) at py4j.Gateway.invoke(Gateway.java:282) at py4j.commands.AbstractCommand.invokeMethod(AbstractCommand.java:132) at py4j.commands.CallCommand.execute(CallCommand.java:79) at py4j.ClientServerConnection.waitForCommands(ClientServerConnection.java:182) at py4j.ClientServerConnection.run(ClientServerConnection.java:106) at java.base/java.lang.Thread.run(Unknown Source) Caused by: java.lang.ClassNotFoundException: org.apache.hadoop.fs.impl.prefetch.PrefetchingStatistics at java.base/jdk.internal.loader.BuiltinClassLoader.loadClass(Unknown Source) at java.base/jdk.internal.loader.ClassLoaders$AppClassLoader.loadClass(Unknown Source) at java.base/java.lang.ClassLoader.loadClass(Unknown Source) ... 35 more
已尝试的解决方案
- 更换多个Spark版本
- 手动添加
hadoop-aws和aws-java-sdk-bundleJar包 - 在Docker外直接运行PySpark/Spark Shell,问题一致
解决方案
这个错误的核心是Hadoop组件版本不兼容:你添加的hadoop-aws-3.3.6.jar依赖Hadoop 3.3.x版本的hadoop-common组件,但apache/spark-py镜像自带的Hadoop版本可能低于3.3.x,导致缺少PrefetchingStatistics这个在Hadoop 3.3.0+才引入的类。
方案1:匹配Spark镜像的Hadoop版本
- 进入容器查看自带Hadoop版本:
hadoop version - 下载与该版本完全一致的
hadoop-awsJar包,替换Dockerfile中的对应文件。比如镜像自带Hadoop 3.2.4,就用hadoop-aws-3.2.4.jar而非3.3.6。
方案2:使用指定Hadoop版本的Spark镜像
直接拉取带对应Hadoop版本的Spark镜像,从根源避免版本冲突:
# 示例:Spark 3.5.0 + Hadoop 3.3.6的镜像 FROM apache/spark-py:v3.5.0-hadoop3.3 COPY ./aws-java-sdk-bundle-1.12.367.jar /opt/spark/jars/aws-java-sdk-bundle-1.12.367.jar COPY . .
注:apache/spark-py的镜像标签格式通常为v<spark-version>-hadoop<hadoop-version>,可按需选择对应版本。
方案3:添加完整Hadoop依赖包(不推荐优先使用)
如果必须用高版本hadoop-aws,需同时替换hadoop-common、hadoop-client等核心Jar包,确保所有Hadoop组件版本一致。但这种方式易引发其他依赖冲突,仅作为备选。
额外配置优化
可添加以下配置避免潜在权限或连接问题:
spark.sparkContext._jsc.hadoopConfiguration().set("fs.s3a.connection.maximum", "100") spark.sparkContext._jsc.hadoopConfiguration().set("fs.s3a.cache.dir", "/tmp/s3a") spark.sparkContext._jsc.hadoopConfiguration().set("fs.s3a.cache.impl", "org.apache.hadoop.fs.s3a.TemporaryAWSCredentialsProvider")
内容的提问来源于stack exchange,提问作者David Sand
相关产品推荐
相关产品推荐

