Databricks Connect运行Spark任务时GraphFrames connectedComponents报错
Databricks Connect调用GraphFrame connectedComponents抛出异常问题
问题描述
通过Databricks Connect运行Spark任务时,调用GraphFrame的connectedComponents方法抛出java.util.NoSuchElementException: None.get异常,但相同代码在Databricks Notebook中可正常运行,且检查checkpoint目录无任何内容。
环境与配置
- 集群版本:10.4 LTS(包含Apache Spark 3.2.1、Scala 2.12)
- SparkSession配置:
spark = ( SparkSession .builder .config( "spark.jars.packages", "graphframes:graphframes:0.8.2-spark3.2-s_2.12" ) .getOrCreate() )
- 业务代码:
edges = genealogy.toDF('src', 'dst') g = GraphFrame(vertices, edges) spark.sparkContext.setCheckpointDir('/') cc = g.connectedComponents()#.cache()
异常信息
py4j.protocol.Py4JJavaError: An error occurred while calling o120.run. : java.util.NoSuchElementException: None.get at scala.None$.get(Option.scala:529) at scala.None$.get(Option.scala:527) at org.apache.spark.sql.util.ProtoSerializer$.fromOption(ProtoSerializer.scala:121) at org.apache.spark.sql.util.ProtoSerializer.deserializeExpr(ProtoSerializer.scala:7393) at org.apache.spark.sql.util.ProtoSerializer.$anonfun$deserializeExpr$5(ProtoSerializer.scala:7411) at scala.collection.TraversableLike.$anonfun$map$1(TraversableLike.scala:286) at scala.collection.Iterator.foreach(Iterator.scala:943) at scala.collection.Iterator.foreach$(Iterator.scala:943) at scala.collection.AbstractIterator.foreach(Iterator.scala:1431) at scala.collection.IterableLike.foreach(IterableLike.scala:74) at scala.collection.IterableLike.foreach$(IterableLike.scala:73) at scala.collection.AbstractIterable.foreach(Iterable.scala:56) at scala.collection.TraversableLike.map(TraversableLike.scala:286) at scala.collection.TraversableLike.map$(TraversableLike.scala:279) at scala.collection.AbstractTraversable.map(Traversable.scala:108) at org.apache.spark.sql.util.ProtoSerializer.deserializeExpr(ProtoSerializer.scala:7411) at org.apache.spark.sql.util.ProtoSerializer.deserializePlan0(ProtoSerializer.scala:4937) at org.apache.spark.sql.util.ProtoSerializer.deserializePlan(ProtoSerializer.scala:4701) at com.databricks.service.SparkServiceRPCHandler.execute0(SparkServiceRPCHandler.scala:666) at com.databricks.service.SparkServiceRPCHandler.$anonfun$executeRPC0$1(SparkServiceRPCHandler.scala:477) at scala.util.DynamicVariable.withValue(DynamicVariable.scala:62) at com.databricks.service.SparkServiceRPCHandler.executeRPC0(SparkServiceRPCHandler.scala:372) at com.databricks.service.SparkServiceRPCHandler$$anon$2.call(SparkServiceRPCHandler.scala:323) at com.databricks.service.SparkServiceRPCHandler$$anon$2.call(SparkServiceRPCHandler.scala:309) at java.util.concurrent.FutureTask.run(FutureTask.java:266) at com.databricks.service.SparkServiceRPCHandler.$anonfun$executeRPC$1(SparkServiceRPCHandler.scala:359) at scala.util.DynamicVariable.withValue(DynamicVariable.scala:62) at com.databricks.service.SparkServiceRPCHandler.executeRPC(SparkServiceRPCHandler.scala:336) at com.databricks.service.SparkServiceRPCServlet.doPost(SparkServiceRPCHandler.scala:167) at javax.servlet.http.HttpServlet.service(HttpServlet.java:523) at javax.servlet.http.HttpServlet.service(HttpServlet.java:590) at org.eclipse.jetty.servlet.ServletHolder.handle(ServletHolder.java:799) at org.eclipse.jetty.servlet.ServletHandler.doHandle(ServletHandler.java:550) at org.eclipse.jetty.server.handler.ScopedHandler.nextScope(ScopedHandler.java:190) at org.eclipse.jetty.servlet.ServletHandler.doScope(ServletHandler.java:501) at org.eclipse.jetty.server.handler.ScopedHandler.handle(ScopedHandler.java:141) at org.eclipse.jetty.server.handler.HandlerWrapper.handle(HandlerWrapper.java:127) at org.eclipse.jetty.server.Server.handle(Server.java:516) at org.eclipse.jetty.server.HttpChannel.lambda$handle$1(HttpChannel.java:388) at org.eclipse.jetty.server.HttpChannel.dispatch(HttpChannel.java:633) at org.eclipse.jetty.server.HttpChannel.handle(HttpChannel.java:380) at org.eclipse.jetty.server.HttpConnection.onFillable(HttpConnection.java:277) at org.eclipse.jetty.io.AbstractConnection$ReadCallback.succeeded(AbstractConnection.java:311) at org.eclipse.jetty.io.FillInterest.fillable(FillInterest.java:105) at org.eclipse.jetty.io.ChannelEndPoint$1.run(ChannelEndPoint.java:104) at org.eclipse.jetty.util.thread.strategy.EatWhatYouKill.runTask(EatWhatYouKill.java:338) at org.eclipse.jetty.util.thread.strategy.EatWhatYouKill.doProduce(EatWhatYouKill.java:315) at org.eclipse.jetty.util.thread.strategy.EatWhatYouKill.tryProduce(EatWhatYouKill.java:173) at org.eclipse.jetty.util.thread.strategy.EatWhatYouKill.run(EatWhatYouKill.java:131) at org.eclipse.jetty.util.thread.ReservedThreadExecutor$ReservedThread.run(ReservedThreadExecutor.java:386) at org.eclipse.jetty.util.thread.QueuedThreadPool.runJob(QueuedThreadPool.java:883) at org.eclipse.jetty.util.thread.QueuedThreadPool$Runner.run(QueuedThreadPool.java:1034) at java.lang.Thread.run(Thread.java:748)
原因分析
- Checkpoint路径错误:设置的
/是客户端本地路径,而非集群可访问的DBFS路径,导致集群无法写入checkpoint数据,触发序列化异常。 - 依赖传递失效:仅在客户端配置
spark.jars.packages,Databricks Connect无法将GraphFrames的jar包正确同步到集群,集群端缺少必要类引发序列化失败。 - 版本不兼容:本地Databricks Connect版本与集群版本不匹配,或GraphFrames版本与Spark/Scala版本存在兼容性冲突,导致执行计划序列化时出现空值。
- 懒加载未触发:
connectedComponents是懒加载操作,仅创建对象未执行,导致序列化逻辑未完整执行。
解决办法
- 修正Checkpoint路径:使用集群可访问的DBFS路径,例如:
确保路径存在且集群有写入权限。spark.sparkContext.setCheckpointDir('/dbfs/tmp/checkpoint') - 统一集群依赖:直接在Databricks集群的「库」页面添加
graphframes:graphframes:0.8.2-spark3.2-s_2.12,或通过集群初始化脚本安装,避免依赖仅存在于客户端。 - 匹配版本:本地Databricks Connect版本必须与集群10.4 LTS完全一致,消除版本不兼容问题。
- 显式触发执行:创建
cc对象后调用Action操作触发任务执行,例如:cc = g.connectedComponents() cc.show() - 简化序列化逻辑:尽量简化GraphFrames的操作逻辑,减少客户端生成的复杂执行计划,降低序列化传递风险。
内容的提问来源于stack exchange,提问作者shammery
相关产品推荐
相关产品推荐

