You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何让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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.19 06:55:22