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库的序列化问题。
方案3:调整Flink ClosureCleaner配置(仅应急)
如果暂时无法重构代码,可关闭Flink的ClosureCleaner序列化检查绕过问题,但这会隐藏潜在的序列化风险,仅建议测试环境临时使用。
val env = StreamExecutionEnvironment.getExecutionEnvironment // 关闭ClosureCleaner的序列化检查 env.getConfig.setClosureCleanerLevel(ClosureCleanerLevel.NONE)
内容的提问来源于stack exchange,提问作者buyalsky
相关产品推荐
相关产品推荐

