Kubernetes Leader Election单元测试失败,请求排查
Kubernetes Fabric8 Leader Election Scala测试失败排查
问题背景
使用Kubernetes Fabric8 Java Client在微服务中实现Leader Election逻辑,参考Fabric8官方代码改编为Scala版本后,单元测试失败。预期能看到Leader当选及卸任的打印信息,但实际leaderLatch.await(10L, SECONDS)返回false,测试抛出TestFailedException。
Scala实现代码
def testAndAssertSingleLeader(client: KubernetesClient, id: String, lock: Lock): Assertion = { // Given val leaderLatch = new CountDownLatch(1) val newLeaderRecord = new AtomicReference[String] val stoppedLeading = new CountDownLatch(1) val executorService = Executors.newSingleThreadExecutor val leaderCallback = new LeaderCallbacks( () => { println("I am the leader now, performing the batch job.") leaderLatch.countDown() }, () => { println("I am no longer the leader.") stoppedLeading.countDown() }, (newLeader: String) => { println(s"A new leader is elected: $newLeader") newLeaderRecord.set(newLeader) }) val leaderElectionConfig = new LeaderElectionConfig( lock, Duration.ofSeconds(15L), Duration.ofSeconds(10L), Duration.ofSeconds(2L), leaderCallback, true, "Integration test leader election configuration" ) // When val leaderElectorTask = executorService.submit( () => assertThrows( classOf[InterruptedException], () => client.leaderElector.withConfig(leaderElectionConfig).build().run() ) ) // Then println(leaderLatch.getCount) assert(leaderLatch.await(10, TimeUnit.SECONDS)) assert(id == newLeaderRecord.get) assert(0 == leaderLatch.getCount) leaderElectorTask.cancel(true) executorService.shutdownNow assert(executorService.awaitTermination(10, TimeUnit.SECONDS)) assert(stoppedLeading.await(10, TimeUnit.SECONDS)) } "K8sLeaderElection" should "create a lease lock" in { // Given server.expect .post .withPath("/apis/coordination.k8s.io/v1/namespaces/namespace/leases") .andReturn(200, null) .once // When - Then testAndAssertSingleLeader(client, "lead-lease", new LeaseLock("namespace", "name", "lead-lease")) }
测试错误日志
Testing started at 22:39 ... Oct 24, 2023 10:39:31 PM okhttp3.mockwebserver.MockWebServer$2 execute INFO: MockWebServer[51585] starting to accept connections 22:39:31.833 [ScalaTest-run] DEBUG io.fabric8.kubernetes.client.utils.HttpClientUtils - Using httpclient io.fabric8.kubernetes.client.okhttp.OkHttpClientFactory factory 22:39:42.108 [ScalaTest-run] DEBUG io.fabric8.kubernetes.client.impl.BaseClient - The client and associated httpclient io.fabric8.kubernetes.client.okhttp.OkHttpClientImpl have been closed, the usage of this or any client using the httpclient will not work after this 1 Oct 24, 2023 10:39:42 PM okhttp3.mockwebserver.MockWebServer$2 acceptConnections INFO: MockWebServer[51585] done accepting connections: Socket closed leaderLatch.await(10L, SECONDS) was false ScalaTestFailureLocation: com.openelectrons.cpo.k8s.K8sLeaderElectionSpec at (K8sLeaderElectionSpec.scala:137) org.scalatest.exceptions.TestFailedException: leaderLatch.await(10L, SECONDS) was false at org.scalatest.Assertions.newAssertionFailedException(Assertions.scala:472) at org.scalatest.Assertions.newAssertionFailedException$(Assertions.scala:471) at org.scalatest.Assertions$.newAssertionFailedException(Assertions.scala:1231) at org.scalatest.Assertions$AssertionsHelper.macroAssert(Assertions.scala:1295) at com.openelectrons.cpo.k8s.K8sLeaderElectionSpec.testAndAssertSingleLeader(K8sLeaderElectionSpec.scala:137) at com.openelectrons.cpo.k8s.K8sLeaderElectionSpec.$anonfun$new$2(K8sLeaderElectionSpec.scala:154)
build.sbt依赖
"org.scalatest" %% "scalatest" % "3.2.9" % Test, "org.mockito" %% "mockito-scala" % "1.16.46" % Test, "org.bouncycastle" % "bcpkix-jdk15on" % "1.68" % Test, "io.fabric8" % "kubernetes-server-mock" % "6.9.0" % Test exclude("com.fasterxml.jackson.core", "jackson-databind")
问题原因及解决方案
1. MockServer返回无效Lease资源
当前MockServer的POST请求返回200和null,但Leader Election逻辑需要有效的Lease对象来确认锁获取成功。返回null会导致客户端无法解析Lease,无法触发Leader当选回调。
修改方案:构造符合K8s API规范的Lease对象返回,确保holderIdentity设置为测试用的leader ID("lead-lease"),同时使用201 Created状态码更贴合K8s API行为:
// Given val lease = new LeaseBuilder() .withMetadata(new ObjectMetaBuilder().withName("name").withNamespace("namespace").build()) .withSpec(new LeaseSpecBuilder() .withHolderIdentity("lead-lease") .withLeaseDurationSeconds(15) .withRenewTime(new Date()) .build()) .build() server.expect .post() .withPath("/apis/coordination.k8s.io/v1/namespaces/namespace/leases") .andReturn(201, lease) .once()
2. MockServer未处理续租请求
Leader Election逻辑会定期更新Lease(续租操作),仅预期一次POST请求不足以覆盖完整流程,需要添加PUT请求的预期:
// 添加续租请求的Mock预期 server.expect .put() .withPath("/apis/coordination.k8s.io/v1/namespaces/namespace/leases/name") .andReturn(200, lease) .times(2) // 根据retryPeriod和测试超时调整次数
3. assertThrows干扰LeaderElector正常运行
在提交的任务中用assertThrows捕获InterruptedException,可能会在Leader回调触发前终止执行流程。应改为捕获异常但不中断逻辑:
val leaderElectorTask = executorService.submit(() => { try { client.leaderElector.withConfig(leaderElectionConfig).build().run() } catch { case _: InterruptedException => // 预期的中断,无需处理 } })
4. 调整超时参数匹配
当前测试await超时为10秒,与LeaderElectionConfig的renewDeadline一致,可能因重试间隔导致回调触发不及时。可以适当延长测试超时时间,或调小retryPeriod:
// 延长await超时到15秒 assert(leaderLatch.await(15, TimeUnit.SECONDS))
内容的提问来源于stack exchange,提问作者joesan
相关产品推荐
相关产品推荐

