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

Spark Structured Streaming部署报错:找不到ByteArraySerializer类

Spark Structured Streaming 部署 Kafka 依赖缺失问题

错误日志

Exception in thread "main" java.lang.reflect.InvocationTargetException
        at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke0(Native Method)
        at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:77)
        at java.base/jdk.internal.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
        at java.base/java.lang.reflect.Method.invoke(Method.java:569)
        at org.apache.spark.deploy.worker.DriverWrapper$.main(DriverWrapper.scala:63)
        at org.apache.spark.deploy.worker.DriverWrapper.main(DriverWrapper.scala)
Caused by: java.lang.NoClassDefFoundError: org/apache/kafka/common/serialization/ByteArraySerializer
        at org.apache.spark.sql.kafka010.KafkaSourceProvider$.<clinit>(KafkaSourceProvider.scala:601)
        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:233)
        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:35)
        at org.apache.spark.sql.streaming.DataStreamReader.loadInternal(DataStreamReader.scala:168)
        at org.apache.spark.sql.streaming.DataStreamReader.load(DataStreamReader.scala:144)
        at DR$.main(DR.scala:112)
        at DR.main(DR.scala)
        ... 6 more
Caused by: java.lang.ClassNotFoundException: org.apache.kafka.common.serialization.ByteArraySerializer
        at java.base/jdk.internal.loader.BuiltinClassLoader.loadClass(BuiltinClassLoader.java:641)
        at java.base/jdk.internal.loader.ClassLoaders$AppClassLoader.loadClass(ClassLoaders.java:188)
        at java.base/java.lang.ClassLoader.loadClass(ClassLoader.java:525)
        ... 17 more
25/04/29 17:19:52 ERROR TransportRequestHandler: Error while invoking RpcHandler#receive() for one-way message.
org.apache.spark.SparkException: Could not find CoarseGrainedScheduler.
        at org.apache.spark.rpc.netty.Dispatcher.postMessage(Dispatcher.scala:178)
        at org.apache.spark.rpc.netty.Dispatcher.postOneWayMessage(Dispatcher.scala:150)
        at org.apache.spark.rpc.netty.NettyRpcHandler.receive(NettyRpcEnv.scala:690)
        at org.apache.spark.network.server.TransportRequestHandler.processOneWayMessage(TransportRequestHandler.java:274)
        at org.apache.spark.network.server.TransportRequestHandler.handle(TransportRequestHandler.java:111)
        at org.apache.spark.network.server.TransportChannelHandler.channelRead0(TransportChannelHandler.java:140)
        at org.apache.spark.network.server.TransportChannelHandler.channelRead0(TransportChannelHandler.java:53)
        at io.netty.channel.SimpleChannelInboundHandler.channelRead(SimpleChannelInboundHandler.java:99)

已尝试的解决方式

  • 使用--packages指定依赖:--packages org.apache.spark:spark-sql-kafka-0-10_2.13:3.4.0
  • 额外添加kafka-clients依赖:--packages org.apache.spark:spark-sql-kafka-0-10_2.13:3.4.0,org.apache.kafka:kafka-clients:3.3.2
  • 手动指定Jar包:--jars spark-streaming-kafka-0-10-assembly_2.13-3.4.0.jar,spark-sql-kafka-0-10_2.13-3.4.0.jar,kafka-clients-3.3.2.jar
  • 去掉kafka-clients的Jar包指定:--jars spark-streaming-kafka-0-10-assembly_2.13-3.4.0.jar,spark-sql-kafka-0-10_2.13-3.4.0.jar
  • 下载所有依赖包到所有Worker节点,通过配置类路径加载:--jars "/jars-kafka-m/*" --conf spark.driver.extraClassPath=/jars-kafka-m/* --conf spark.executor.extraClassPath=/jars-kafka-m/*
  • 将kafka-clients-3.3.2打包进应用Jar包

环境信息

  • Spark版本:3.4.0
  • 应用Scala版本:2.13.13
  • Sbt版本:1.9.9
  • Kafka Broker版本:3.4.0

疑问

通过--packages获取的kafka-client版本是3.3.2,但Kafka Broker的libs文件夹中是kafka-clients-3.4.0,不确定这是否是报错原因,寻求解决指导。


解决方案

1. 版本兼容说明

Spark 3.4.0官方依赖的kafka-clients版本是3.3.2,这个版本和Kafka Broker 3.4.0完全兼容(Kafka客户端与Broker跨小版本兼容),版本差异不是报错原因。问题核心是依赖未正确分发到Driver和Executor节点。

2. 优先使用--packages参数

--packages是Spark管理外部依赖最可靠的方式,会自动下载依赖并分发到所有集群节点,无需手动处理Jar包。执行命令时确保集群节点能访问Maven中央仓库,若使用内部仓库,可通过--conf spark.jars.repositories=http://your-repo-url指定。

完整提交命令示例:

spark-submit \
  --class DR \
  --master spark://<master-ip>:7077 \
  --packages org.apache.spark:spark-sql-kafka-0-10_2.13:3.4.0 \
  your-application.jar

3. 手动Jar包部署注意事项

若必须手动上传Jar包:

  • 所有Worker节点的Jar包路径需完全一致,且Spark运行用户有读取权限
  • 避免同时使用--jars和extraClassPath,易引发类加载冲突
  • 不要将kafka-clients打包进应用Jar,会与Spark分发的依赖冲突

4. CoarseGrainedScheduler错误排查

该错误是Driver启动失败的连锁问题,解决依赖缺失后会自动消失。若依赖问题解决后仍出现,检查:

  • Master节点是否正常运行
  • Driver节点能正常连接Master的7077端口
  • Spark集群资源配置足够启动Driver

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 07:20:05