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

Ubuntu 22.04下Spark Streaming报错:找不到ByteArraySerializer类

问题:Ubuntu 22.04上Spark Streaming应用找不到Kafka相关类

我在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支持。

解决步骤:

  1. 下载对应版本的Kafka客户端Jar:
    Spark 3.3.0官方依赖的Kafka客户端版本是2.8.1,需要下载kafka-clients-2.8.1.jar(源码包可选)。

  2. 放置Jar到Spark目录:
    将下载好的kafka-clients-2.8.1.jar放到/usr/local/spark/jars目录,确保文件权限为Spark进程可读取。

  3. 清理依赖缓存(可选):
    如果之前通过spark.jars.packages拉取过依赖,可能存在缓存异常,执行以下命令删除缓存:

    rm -rf ~/.ivy2/cache/org.apache.kafka
    
  4. 重新运行应用:
    可以移除代码中的spark.jars.packages配置(已手动放置Jar),或保留让Spark自动处理依赖,手动放置Jar的方式更稳定。

额外说明:

Windows上能正常运行是因为Spark通过spark.jars.packages自动下载了所有依赖(包括Kafka客户端Jar),而Ubuntu下可能因网络或缓存问题未拉全依赖,导致缺失底层类。

内容的提问来源于stack exchange,提问作者Ivan

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.22 13:57:17