Spark Connect自定义拦截器加载失败:类型转换异常求助
问题:Spark Connect自定义拦截器启动时出现类型转换异常
我正尝试按照Spark官方文档实现自定义Spark Connect拦截器,编写了如下简单拦截器代码:
package interceptorserver; import io.grpc.Metadata; import io.grpc.ServerCall; import io.grpc.ServerCall.Listener; import io.grpc.ServerCallHandler; import io.grpc.ServerInterceptor; public class Interceptor implements ServerInterceptor{ @Override public <ReqT, RespT> Listener<ReqT> interceptCall(ServerCall<ReqT, RespT> call, Metadata headers, ServerCallHandler<ReqT, RespT> next) { System.out.println("Hello world"); return next.startCall(call, headers); } }
编译后通过以下命令启动Spark Connect服务器:
./start-connect-server.sh \ --packages org.apache.spark:spark-connect_2.12:3.4.1 \ --jars Interceptor.jar \ --conf spark.connect.grpc.interceptor.classes=interceptorserver.Interceptor
但启动时出现如下错误:
23/07/29 01:17:00 ERROR SparkConnectServer: Error starting Spark Connect server org.apache.spark.SparkException: [CONNECT.INTERCEPTOR_RUNTIME_ERROR] Generic Spark Connect error. Error instantiating GRPC interceptor: class interceptorserver.Interceptor cannot be cast to class org.sparkproject.connect.grpc.ServerInterceptor (interceptorserver.Interceptor and org.sparkproject.connect.grpc.ServerInterceptor are in unnamed module of loader org.apache.spark.util.MutableURLClassLoader @a5272be) at org.apache.spark.sql.connect.service.SparkConnectInterceptorRegistry$.createInstance(SparkConnectInterceptorRegistry.scala:99) at org.apache.spark.sql.connect.service.SparkConnectInterceptorRegistry$.$anonfun$createConfiguredInterceptors$4(SparkConnectInterceptorRegistry.scala:67) at scala.collection.TraversableLike.$anonfun$map$1(TraversableLike.scala:286) at scala.collection.IndexedSeqOptimized.foreach(IndexedSeqOptimized.scala:36) at scala.collection.IndexedSeqOptimized.foreach$(IndexedSeqOptimized.scala:33) at scala.collection.mutable.ArrayOps$ofRef.foreach(ArrayOps.scala:198) ...
我最初以为org.sparkproject.connect.grpc.ServerInterceptor与io.grpc.ServerInterceptor是不同接口,但查看Spark源码及官方文档后确认,Spark确实使用io.grpc.ServerInterceptor。随后我编写测试验证自定义拦截器确实实现了该接口:
/* * This Java source file was generated by the Gradle 'init' task. */ package interceptorserver; import org.junit.Test; import org.junit.Assert; public class LibraryTest { @Test public void someLibraryMethodReturnsTrue() { Interceptor classUnderTest = new Interceptor(); Assert.assertTrue(classUnderTest instanceof io.grpc.ServerInterceptor); } }
测试通过,但问题仍存在。请问我哪里操作有误?为何会出现类型转换异常?
内容的提问来源于stack exchange,提问作者Matheus
相关产品推荐
相关产品推荐

