Ubuntu 22.04下Spark Streaming报错:找不到ByteArraySerializer类
我在Ubuntu 22.04上运行Spark Streaming应用时遇到错误,相同配置在Windows上能正常运行,但Ubuntu上提示找不到Jar文件中的类。
SparkSession初始化配置
spark = SparkSession \ .builder \ .appName("File Streaming PostgreSQL") \ .master("local[3]") \ .config("spark.streaming.stopGracefullyOnShutdown", "true") \ .config("spark.jars.packages", "org.apache.spark:spark-avro_2.12:3.3.0,org.apache.spark:spark-sql-kafka-0-10_2.12:3.3.0") \ .config("spark.sql.shuffle.partitions", 2) \ .getOrCreate()
已放置的Jar文件
我已将以下Avro和Kafka相关Jar文件放到/usr/local/spark/jars目录:
- spark-sql-kafka-0-10_2.12-3.3.0.jar
- spark-sql-kafka-0-10_2.12-3.3.0-tests.jar
- spark-sql-kafka-0-10_2.12-3.3.0-javadoc.jar
- spark-sql-kafka-0-10_2.12-3.3.0-sources.jar
- spark-sql-kafka-0-10_2.12-3.3.0-test-sources.jar
- spark-avro_2.12-3.3.0.jar
- spark-avro_2.12-3.3.0-tests.jar
- spark-avro_2.12-3.3.0-javadoc.jar
- spark-avro_2.12-3.3.0-sources.jar
- spark-avro_2.12-3.3.0-test-sources.jar
环境版本
- Spark 3.3.0
- Scala 2.12.14
- OpenJDK 64-Bit Server VM(Java 11.0.16)
错误日志
File "/usr/local/spark/python/lib/py4j-0.10.9.5-src.zip/py4j/protocol.py", line 326, in get_return_value py4j.protocol.Py4JJavaError: An error occurred while calling o39.load. : java.lang.NoClassDefFoundError: org/apache/kafka/common/serialization/ByteArraySerializer at org.apache.spark.sql.kafka010.KafkaSourceProvider$.<init>(KafkaSourceProvider.scala:601) at org.apache.spark.sql.kafka010.KafkaSourceProvider$.<clinit>(KafkaSourceProvider.scala) at org.apache.spark.sql.kafka010.KafkaSourceProvider.org$apache$spark$sql$kafka010$KafkaSourceProvider$$validateStreamOptions(KafkaSourceProvider.scala:338) at org.apache.spark.sql.kafka010.KafkaSourceProvider.sourceSchema(KafkaSourceProvider.scala:71) at org.apache.spark.sql.execution.datasources.DataSource.sourceSchema(DataSource.scala:236) at org.apache.spark.sql.execution.datasources.DataSource.sourceInfo$lzycompute(DataSource.scala:118) at org.apache.spark.sql.execution.datasources.DataSource.sourceInfo(DataSource.scala:118) at org.apache.spark.sql.execution.streaming.StreamingRelation$.apply(StreamingRelation.scala:34) at org.apache.spark.sql.streaming.DataStreamReader.loadInternal(DataStreamReader.scala:168) at org.apache.spark.sql.streaming.DataStreamReader.load(DataStreamReader.scala:144) at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke0(Native Method) at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62) at java.base/jdk.internal.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43) at java.base/java.lang.reflect.Method.invoke(Method.java:566) at py4j.reflection.MethodInvoker.invoke(MethodInvoker.java:244) at py4j.reflection.ReflectionEngine.invoke(ReflectionEngine.java:357) 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(Thread.java:829) Caused by: java.lang.ClassNotFoundException: org.apache.kafka.common.serialization.ByteArraySerializer at java.base/jdk.internal.loader.BuiltinClassLoader.loadClass(BuiltinClassLoader.java:581) at java.base/jdk.internal.loader.ClassLoaders$AppClassLoader.loadClass(ClassLoaders.java:178) at java.base/java.lang.ClassLoader.loadClass(ClassLoader.java:522) ... 22 more
已配置的环境变量(.bashrc)
#configuration for local Spark and Hadoop SPARK_HOME=/usr/local/spark-3.3.0-bin-hadoop3 export PATH=$PATH:$SPARK_HOME/bin:$SPARK_HOME/sbin export PYSPARK_PYTHON=/usr/bin/python3 export JAVA_HOME=/usr/lib/jvm/java-11-openjdk-amd64 export PATH=$PATH:$JAVA_HOME/bin
补充:该应用在Windows上正常运行,Ubuntu上其他不涉及SQL和Avro序列化的Spark应用也能正常运行。
这个错误是因为缺少Kafka客户端的核心Jar包,你只放了Spark封装的Kafka相关Jar,但Spark的Kafka依赖需要底层的Kafka客户端Jar支持。
解决步骤:
下载对应版本的Kafka客户端Jar:
Spark 3.3.0官方依赖的Kafka客户端版本是2.8.1,需要下载kafka-clients-2.8.1.jar(源码包可选)。放置Jar到Spark目录:
将下载好的kafka-clients-2.8.1.jar放到/usr/local/spark/jars目录,确保文件权限为Spark进程可读取。清理依赖缓存(可选):
如果之前通过spark.jars.packages拉取过依赖,可能存在缓存异常,执行以下命令删除缓存:rm -rf ~/.ivy2/cache/org.apache.kafka重新运行应用:
可以移除代码中的spark.jars.packages配置(已手动放置Jar),或保留让Spark自动处理依赖,手动放置Jar的方式更稳定。
额外说明:
Windows上能正常运行是因为Spark通过spark.jars.packages自动下载了所有依赖(包括Kafka客户端Jar),而Ubuntu下可能因网络或缓存问题未拉全依赖,导致缺失底层类。
内容的提问来源于stack exchange,提问作者Ivan

