如何连接远程Spark集群执行作业?解决类找不到异常
外部应用连接Bitnami Spark集群执行作业的问题与解决
问题背景
通过bitnami/spark Helm Chart部署了Spark集群,包含1个Master节点和2个Worker节点,希望外部应用连接集群执行作业,但遇到ClassNotFoundException错误。
作业代码
package com.sompackage; import org.apache.spark.SparkContext; import org.apache.spark.rdd.RDD; import org.apache.spark.sql.SparkSession; import java.util.*; import java.util.stream.Collectors; import java.util.stream.IntStream; import scala.collection.JavaConverters; public class SparkTestMain { public static void main(String[] args) { SparkSession spark = SparkSession.builder().master("sc://master-node-url:masterport") .appName("SparkPi") .config("spark.submit.deployMode","client") .config("spark.driver.host", "my-pub-ip") .config("spark.driver.bindAddress", "my-private-ip") .config("spark.driver.port", "1080") .config("spark.fileserver.port", "1081") .config("spark.broadcast.port", "1082") .config("spark.replClassServer.port", "1083") .config("spark.blockManager.port", "1084") .config("spark.executor.port", "1085") .getOrCreate(); SparkContext sc = spark.sparkContext(); int slices = (args.length == 1) ? Integer.parseInt(args[0]) : 2; int n = 100 * slices; List<Integer> l = new ArrayList<>(n); for (int i = 0; i < n; i++) { l.add(i); } List<Integer> seqNumList = IntStream.rangeClosed(10, 20).boxed().collect(Collectors.toList()); RDD<Integer> numRDD = sc .parallelize(JavaConverters.asScalaIteratorConverter(seqNumList.iterator()).asScala() .toSeq(), 2, scala.reflect.ClassTag$.MODULE$.apply(Integer.class)); numRDD.toJavaRDD().foreach(x -> System.out.println(x)); spark.stop(); } }
报错信息
java.lang.ClassNotFoundException: com.somepackage.SparkTestMain at java.base/java.net.URLClassLoader.findClass(URLClassLoader.java:445) at java.base/java.lang.ClassLoader.loadClass(ClassLoader.java:592) at java.base/java.lang.ClassLoader.loadClass(ClassLoader.java:525) at java.base/java.lang.Class.forName0(Native Method) at java.base/java.lang.Class.forName(Class.java:467) at org.apache.spark.serializer.JavaDeserializationStream$$anon$1.resolveClass(JavaSerializer.scala:71) at java.base/java.io.ObjectInputStream.readNonProxyDesc(ObjectInputStream.java:2034) at java.base/java.io.ObjectInputStream.readClassDesc(ObjectInputStream.java:1898) at java.base/java.io.ObjectInputStream.readClass(ObjectInputStream.java:1871) at java.base/java.io.ObjectInputStream.readObject0(ObjectInputStream.java:1704) at java.base/java.io.ObjectInputStream$FieldValues.<init>(ObjectInputStream.java:2606) at java.base/java.io.ObjectInputStream.readSerialData(ObjectInputStream.java:2457) at java.base/java.io.ObjectInputStream.readOrdinaryObject(ObjectInputStream.java:2257) at java.base/java.io.ObjectInputStream.readObject0(ObjectInputStream.java:1733) at java.base/java.io.ObjectInputStream.readArray(ObjectInputStream.java:2157) at java.base/java.io.ObjectInputStream.readObject0(ObjectInputStream.java:1721) at java.base/java.io.ObjectInputStream$FieldValues.<init>(ObjectInputStream.java:2606) at java.base/java.io.ObjectInputStream.readSerialData(ObjectInputStream.java:2457) at java.base/java.io.ObjectInputStream.readOrdinaryObject(ObjectInputStream.java:2257) at java.base/java.io.ObjectInputStream.readObject0(ObjectInputStream.java:1733) at java.base/java.io.ObjectInputStream.readArray(ObjectInputStream.java:2157) at java.base/java.io.ObjectInputStream.readObject0(ObjectInputStream.java:1721) at java.base/java.io.ObjectInputStream$FieldValues.<init>(ObjectInputStream.java:2606) at java.base/java.io.ObjectInputStream.readSerialData(ObjectInputStream.java:2457) at java.base/java.io.ObjectInputStream.readOrdinaryObject(ObjectInputStream.java:2257) at java.base/java.io.ObjectInputStream.readObject0(ObjectInputStream.java:1733) at java.base/java.io.ObjectInputStream.readArray(ObjectInputStream.java:2157) at java.base/java.io.ObjectInputStream.readObject0(ObjectInputStream.java:1721) at java.base/java.io.ObjectInputStream$FieldValues.<init>(ObjectInputStream.java:2606) at java.base/java.io.ObjectInputStream.readSerialData(ObjectInputStream.java:2457) at java.base/java.io.ObjectInputStream.readOrdinaryObject(ObjectInputStream.java:2257) at java.base/java.io.ObjectInputStream.readObject0(ObjectInputStream.java:1733) at java.base/java.io.ObjectInputStream$FieldValues.<init>(ObjectInputStream.java:2606) at java.base/java.io.ObjectInputStream.readSerialData(ObjectInputStream.java:2457) at java.base/java.io.ObjectInputStream.readOrdinaryObject(ObjectInputStream.java:2257) at java.base/java.io.ObjectInputStream.readObject0(ObjectInputStream.java:1733) at java.base/java.io.ObjectInputStream.readObject(ObjectInputStream.java:509) at java.base/java.io.ObjectInputStream.readObject(ObjectInputStream.java:467) at org.apache.spark.serializer.JavaDeserializationStream.readObject(JavaSerializer.scala:87) at org.apache.spark.serializer.JavaSerializerInstance.deserialize(JavaSerializer.scala:129) at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:86) at org.apache.spark.TaskContext.runTaskWithListeners(TaskContext.scala:161) at org.apache.spark.scheduler.Task.run(Task.scala:141) at org.apache.spark.executor.Executor$TaskRunner.$anonfun$run$4(Executor.scala:620) at org.apache.spark.util.SparkErrorUtils.tryWithSafeFinally(SparkErrorUtils.scala:64) at org.apache.spark.util.SparkErrorUtils.tryWithSafeFinally$(SparkErrorUtils.scala:61) at org.apache.spark.util.Utils$.tryWithSafeFinally(Utils.scala:94) at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:623) at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1136) at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:635) at java.base/java.lang.Thread.run(Thread.java:840)
已尝试操作
更换Java版本和不同Spark版本,问题仍未解决。
解决方案
1. 确保作业Jar包被Worker节点加载
- 编译包含
SparkTestMain类的Jar包,在SparkSession配置中添加spark.jars参数指定Jar包路径:
若Jar包不在本地,可上传至集群共享存储(如HDFS),路径改为SparkSession spark = SparkSession.builder() .master("spark://master-node-url:7077") .appName("SparkPi") .config("spark.jars", "file:///local/path/to/your/job.jar") // 其他配置... .getOrCreate();hdfs:///path/to/job.jar。
2. 修正包名拼写错误
代码中包声明为com.sompackage,但报错提示找不到com.somepackage(多了一个'e'),需统一所有地方的包名拼写,确保编译后的类路径与代码声明一致。
3. 修正Master地址协议与端口
Bitnami Spark集群的Master默认使用spark://协议,端口为7077,需将代码中的sc://master-node-url:masterport改为spark://<master-service-name>:7077,其中<master-service-name>是Kubernetes中Spark Master的服务名称。
4. 验证网络连通性
- 确保外部机器与Spark集群网络互通,开放代码中配置的端口(1080-1085),Worker节点能主动访问
spark.driver.host指定的公网IP和对应端口。 - 若使用Kubernetes集群,需为Driver的端口配置NodePort或LoadBalancer类型的Service,保证Worker能访问到Driver。
5. 确认Deploy Mode适配场景
client模式下,Driver运行在外部应用所在机器,需保证Worker能访问到Driver的IP和端口;- 若集群无法直接访问外部机器,建议切换为
cluster模式,Driver将运行在集群内部,此时需配置spark.driver.host为集群内部可访问的地址。
内容的提问来源于stack exchange,提问作者Ahsan A. Ishan
相关产品推荐
相关产品推荐

