You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

使用PySpark读取Minio数据时触发NoClassDefFoundError求助

PySpark读取Minio时出现NoClassDefFoundError:org/apache/hadoop/fs/impl/prefetch/PrefetchingStatistics

环境

  • 设备:Macbook Air M1
  • Minio:本地9000端口运行
  • Spark:使用apache/spark-py官方镜像

复现步骤

  1. 本地启动Minio实例(端口9000)
  2. 使用以下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 . .
  1. 运行容器:
docker run --network="host" -it <image-name> /opt/spark/bin/pyspark
  1. 在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-bundle Jar包
  • 在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版本

  1. 进入容器查看自带Hadoop版本:
    hadoop version
    
  2. 下载与该版本完全一致的hadoop-aws Jar包,替换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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.13 17:45:03