如何让java.util.function.Predicate序列化?Resilience4j重试配置序列化报错
问题:Resilience4j重试配置在Spark中触发序列化异常
在Spark的MapElementsExec执行阶段,使用Resilience4j配置重试逻辑时抛出java.io.NotSerializableException。核心原因是RetryConfig内部的exceptionPredicate字段为Lambda实现的Predicate实例,而Lambda默认不支持序列化。尽管自定义的GetTaxResilience组件已实现Serializable,但依赖的Resilience4j组件内部存在不可序列化的Lambda,导致整体序列化失败。
报错堆栈信息
App > Serialization stack: App > - object not serializable (class: java.util.function.Predicate$$Lambda$309/931003277, value: java.util.function.Predicate$$Lambda$309/931003277@5c97826a) App > - field (class: io.github.resilience4j.retry.RetryConfig, name: exceptionPredicate, type: interface java.util.function.Predicate) App > - object (class io.github.resilience4j.retry.RetryConfig, io.github.resilience4j.retry.RetryConfig@20496d9f) App > - field (class: determination.resilience.GetTaxResilience, name: retryConfig, type: class io.github.resilience4j.retry.RetryConfig) App > - object (class determination.resilience.GetTaxResilience, determination.resilience.GetTaxResilience@af49e79) App > - field (class: vertex.VertexTaxComplianceOperator, name: getTaxResilience, type: class determination.resilience.GetTaxResilience) App > - object (class vertex.VertexTaxComplianceOperator, vertex.VertexTaxComplianceOperator@3e68a323) App > - field (class: vertex.VertexTaxComplianceOperator$$anonfun$5, name: $outer, type: class vertex.VertexTaxComplianceOperator) App > - object (class vertex.VertexTaxComplianceOperator$$anonfun$5, <function1>) App > - field (class: org.apache.spark.sql.execution.MapElementsExec, name: func, type: class java.lang.Object) App > - object (class org.apache.spark.sql.execution.MapElementsExec, MapElements <function1>, obj#512:
相关代码片段
自定义GetTaxResilience组件
@Component class GetTaxResilience extends Serializable { private val logger = LoggerFactory.getLogger(this.getClass) val retryConfig = RetryConfig.custom() .maxAttempts(RESILIENCE_RETRY_MAX_ATTEMPTS) .retryExceptions(classOf[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 = { // default retry instance configured in yaml logger.info("Retrieving default retry instance: {}", GET_VERTEX_TAX_RETRY_DEFAULT_NAME) retryInstance } }
调用代码
@Autowired var getTaxResilience: GetTaxResilience = null withRetry(getTaxResilience.getDefaultRetryInstance())
Spark相关代码
val calculateResultDF = preparedDataSetDF.mapPartitions(iterator => { val calculateResultDF = iterator.map(row => { //retry stuff }) calculateResultDF }).toDF() spark.sql(sqlBuilder.createStagingTableInsertIntoSQL(jobInfoCase.run_id)).toDF() insertIntoExceptionArchiveTable(spark, jobInfoCase.run_id)
解决办法
1. 替换Lambda为可序列化的匿名内部类
Lambda默认不实现Serializable,因此将retryExceptions替换为retryPredicate,传入实现Serializable的匿名内部类:
val retryConfig = RetryConfig.custom() .maxAttempts(RESILIENCE_RETRY_MAX_ATTEMPTS) .retryPredicate(new Predicate[Throwable] with Serializable { override def test(t: Throwable): Boolean = { classOf[StatusRuntimeException].isInstance(t) } }) .intervalFunction(IntervalFunction.ofRandomized(RESILIENCE_RETRY_INTERVAL_MILLISECONDS)) .build
这样exceptionPredicate字段就变成了可序列化的实例,避免序列化失败。
2. 延迟初始化Retry相关对象
将retryConfig、retryRegistry、retryInstance改为lazy val,延迟初始化直到Executor端第一次调用时创建,避免把Driver端的不可序列化对象传递到Executor:
@Component class GetTaxResilience extends Serializable { private val logger = LoggerFactory.getLogger(this.getClass) lazy val retryConfig = RetryConfig.custom() .maxAttempts(RESILIENCE_RETRY_MAX_ATTEMPTS) .retryPredicate(new Predicate[Throwable] with Serializable { override def test(t: Throwable): Boolean = { classOf[StatusRuntimeException].isInstance(t) } }) .intervalFunction(IntervalFunction.ofRandomized(RESILIENCE_RETRY_INTERVAL_MILLISECONDS)) .build lazy val retryRegistry = RetryRegistry.of(retryConfig) 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 } }
3. 在Spark算子内部创建Retry实例
直接在mapPartitions内部初始化Retry实例,完全避免序列化Resilience4j相关对象:
val calculateResultDF = preparedDataSetDF.mapPartitions(iterator => { // 每个分区独立初始化Retry实例 val retryConfig = RetryConfig.custom() .maxAttempts(RESILIENCE_RETRY_MAX_ATTEMPTS) .retryExceptions(classOf[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) iterator.map(row => { // 执行带重试的业务逻辑 val result = Retry.decorateSupplier(retryInstance, () => { // 你的业务逻辑代码 }).get() result }) }).toDF()
内容的提问来源于stack exchange,提问作者Nick
相关产品推荐
相关产品推荐

