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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.07 21:15:01