实现Serializable的类调用函数触发NotSerializableException如何解决?
问题分析
问题出在你通过asJavaPredicate(exception => exception.isInstanceOf[StatusRuntimeException])创建的Lambda表达式上。虽然RetryConfig本身实现了Serializable,但这个Lambda并未实现java.io.Serializable接口。当Spark需要序列化GetTaxResilience实例(该实例被mainClass持有,而Spark任务会将相关对象序列化分发到Executor节点),就会触发这个序列化异常。
解决方案
方案1:移除重复的重试配置(最简单高效)
你同时调用了retryExceptions(classOf[StatusRuntimeException])和retryOnException(...),两者逻辑完全一致——都是仅对StatusRuntimeException进行重试。保留retryExceptions即可,它内部会自动创建可序列化的Predicate实现,无需手动传入Lambda:
class GetTaxResilience extends Serializable { private val logger = LoggerFactory.getLogger(this.getClass) val retryConfig = RetryConfig.custom() .maxAttempts(RESILIENCE_RETRY_MAX_ATTEMPTS) .retryExceptions(classOf[StatusRuntimeException]) // 移除这一行冗余配置 // .retryOnException(asJavaPredicate(exception => exception.isInstanceOf[StatusRuntimeException])) .intervalFunction(IntervalFunction.ofRandomized(RESILIENCE_RETRY_INTERVAL_MILLISECONDS)) .build val retryRegistry = RetryRegistry.of(retryConfig) val retryInstance = retryRegistry.retry(GET_VERTEX_TAX_RETRY_DEFAULT_NAME, retryConfig) def getDefaultRetryInstance(): Retry = { logger.info("Retrieving default retry instance: {}", GET_VERTEX_TAX_RETRY_DEFAULT_NAME) retryInstance } }
方案2:使用可序列化的Predicate替代Lambda
如果确实需要自定义Predicate逻辑(无法用retryExceptions覆盖),不要用Lambda,改用实现Serializable的匿名内部类或自定义类:
方式A:匿名内部类
.retryOnException(new java.util.function.Predicate[Throwable] with Serializable { override def test(t: Throwable): Boolean = { t.isInstanceOf[StatusRuntimeException] } })
方式B:自定义可序列化Predicate类
// 单独定义可序列化的Predicate类 class StatusRuntimeExceptionPredicate extends java.util.function.Predicate[Throwable] with Serializable { override def test(t: Throwable): Boolean = { t.isInstanceOf[StatusRuntimeException] } } // 在RetryConfig中使用 val retryConfig = RetryConfig.custom() // ...其他配置 .retryOnException(new StatusRuntimeExceptionPredicate()) .build
方案3:标记非序列化字段为transient(适配Driver端专属场景)
如果retryConfig、retryRegistry等对象仅在Driver端使用,不需要序列化到Executor,可以标记为transient,并结合lazy val延迟初始化,避免反序列化后出现空指针:
class GetTaxResilience extends Serializable { private val logger = LoggerFactory.getLogger(this.getClass) // 标记为transient避免序列化,用lazy val保证首次调用时初始化 @transient private lazy val retryConfig = RetryConfig.custom() .maxAttempts(RESILIENCE_RETRY_MAX_ATTEMPTS) .retryExceptions(classOf[StatusRuntimeException]) .retryOnException(asJavaPredicate(exception => exception.isInstanceOf[StatusRuntimeException])) .intervalFunction(IntervalFunction.ofRandomized(RESILIENCE_RETRY_INTERVAL_MILLISECONDS)) .build @transient private lazy val retryRegistry = RetryRegistry.of(retryConfig) @transient private lazy val retryInstance = retryRegistry.retry(GET_VERTEX_TAX_RETRY_DEFAULT_NAME, retryConfig) def getDefaultRetryInstance(): Retry = { logger.info("Retrieving default retry instance: {}", GET_VERTEX_TAX_RETRY_DEFAULT_NAME) retryInstance } }
内容的提问来源于stack exchange,提问作者Nick
相关产品推荐
相关产品推荐

