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

Flink Scala单元测试中Mock CaffeineHelper序列化失败求助

解决Scala Flink单元测试中CaffeineHelper Mock序列化问题

方案1:手动实现测试替身(Test Double)替代Mock库

直接编写一个简单的CaffeineHelper实现类,完全控制其序列化行为,避免Mock库带来的非序列化问题。

// 测试用的CaffeineHelper实现,基于内存Map,天然支持序列化
class TestCaffeineHelper[V](initialData: Map[String, V] = Map.empty) extends CaffeineHelper[V] {
  // 测试环境用普通Map替代Caffeine缓存
  private val mockCache = scala.collection.mutable.Map[String, V](initialData.toSeq: _*)
  
  // 覆盖父类的transient cache(测试中不会用到实际Caffeine实例)
  @transient override protected var cache: CaffeineCache[V] = null

  override def get(id: String): Option[V] = mockCache.get(id)

  override def put(key: String)(value: V): Unit = mockCache.put(key, value)
}

测试中使用该类:

"EnrichWithUserMetaMapper" should "enrich with brand" in {
  val env = StreamExecutionEnvironment.getExecutionEnvironment

  // 预存测试数据的测试替身
  val testCache = new TestCaffeineHelper[Long](Map("user1" -> 4L))

  val ageMapper = new EnrichWithUserMeta() {
    override def open(parameters: Configuration): Unit = {
      userMetadataCache = testCache
    }
  }

  val stream = env.fromElements(User("user1"))

  val enrichedProducts = AsyncDataStream
    .unorderedWait(stream, ageMapper, 5, TimeUnit.SECONDS, 1)
    .executeAndCollect(1)

  enrichedProducts.head shouldBe UserWithMeta("user1", 4L)
}

原理:测试替身完全由自定义实现,没有引入Mock库的非序列化内部对象,天然符合Serializable要求,Flink可正常序列化。

方案2:重构代码,通过构造函数注入CaffeineHelper

修改EnrichWithUserMeta,将CaffeineHelper作为构造参数传入,而非在open方法中初始化,测试时直接传入测试替身,无需重写方法。

重构后的EnrichWithUserMeta:

class EnrichWithUserMeta(
  protected val userMetadataCache: CaffeineHelper[Long]
) extends RichAsyncFunction[User, UserWithMeta] {
  // 无参构造函数,兼容生产环境原有逻辑
  def this() = this(CaffeineHelper.newCache[Long](10.seconds))

  override def open(parameters: Configuration): Unit = {
    super.open(parameters)
  }

  override def asyncInvoke(input: User, resultFuture: ResultFuture[UserWithMeta]): Unit = {
    val age = getUserAge(input.id)
    resultFuture.complete(Seq(UserWithMeta(input.id, age)))
  }

  private def getUserAge(userId: String): Long = {
    userMetadataCache.get(userId) match {
      case Some(age) => age
      case None => {
        val age = getAgeFromDB(userId)
        userMetadataCache.put(userId)(age)
        age
      }
    }
  }

  // 提取DB查询方法为protected,方便测试时重写
  protected def getAgeFromDB(userId: String): Long = {
    // 生产环境DB查询逻辑
    ???
  }
}

测试代码:

"EnrichWithUserMetaMapper" should "enrich with brand" in {
  val env = StreamExecutionEnvironment.getExecutionEnvironment

  val testCache = new TestCaffeineHelper[Long](Map("user1" -> 4L))
  // 直接通过构造函数注入测试替身
  val ageMapper = new EnrichWithUserMeta(testCache) {
    // 可选:重写DB查询方法模拟数据
    override protected def getAgeFromDB(userId: String): Long = 4L
  }

  val stream = env.fromElements(User("user1"))

  val enrichedProducts = AsyncDataStream
    .unorderedWait(stream, ageMapper, 5, TimeUnit.SECONDS, 1)
    .executeAndCollect(1)

  enrichedProducts.head shouldBe UserWithMeta("user1", 4L)
}

原理:依赖注入让测试逻辑更清晰,无需修改open方法,直接传入可控的测试替身,彻底避免Mock库的序列化问题。

如果暂时无法重构代码,可关闭Flink的ClosureCleaner序列化检查绕过问题,但这会隐藏潜在的序列化风险,仅建议测试环境临时使用。

val env = StreamExecutionEnvironment.getExecutionEnvironment
// 关闭ClosureCleaner的序列化检查
env.getConfig.setClosureCleanerLevel(ClosureCleanerLevel.NONE)

内容的提问来源于stack exchange,提问作者buyalsky

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.16 17:49:49