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
相关产品推荐
相关产品推荐

