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

Fabric8 JUnit测试监听自定义资源无WatchEvent触发问题

Fabric8 自定义资源Watcher测试无事件触发问题修复

问题根因

  • JUnit版本注解混用
    代码中用JUnit4的@Rule注解启动KubernetesServer,但测试方法、断言用的是JUnit5的注解。JUnit5不识别@Rule,Mock Server不会完成初始化,预设的响应规则完全不生效。
  • 依赖版本冲突
    引入的io.fabric8:kubernetes-api:3.0.12是过老的不兼容版本,和5.9.0的kubernetes-client存在核心类冲突,会导致Watch事件解析、WebSocket连接逻辑异常。5.x版本客户端已经内置所有需要的API类,直接删除该依赖即可。
  • CRD类型注册时机错误
    KubernetesDeserializer.registerCustomKind调用放在Mock规则设置、客户端初始化之后,Mock Server无法识别自定义资源类型,无法序列化发出对应事件。
  • Watch路径匹配规则失效
    Fabric8 5.x发起Watch请求时会自动拼接resourceVersion查询参数,原代码写的精确路径只能匹配无额外参数的请求,实际请求带了额外参数会导致规则匹配失败,不会返回预设的WebSocket流。
  • 冗余依赖:引入了VertxExtension但全程未使用,可直接移除避免不必要的类加载问题。

修复后可运行代码

import io.fabric8.kubernetes.api.model.Condition;
import io.fabric8.kubernetes.api.model.KubernetesResourceList;
import io.fabric8.kubernetes.api.model.ObjectMetaBuilder;
import io.fabric8.kubernetes.api.model.WatchEvent;
import io.fabric8.kubernetes.api.model.WatchEventBuilder;
import io.fabric8.kubernetes.client.CustomResource;
import io.fabric8.kubernetes.client.KubernetesClient;
import io.fabric8.kubernetes.client.Watch;
import io.fabric8.kubernetes.client.Watcher;
import io.fabric8.kubernetes.client.WatcherException;
import io.fabric8.kubernetes.client.dsl.MixedOperation;
import io.fabric8.kubernetes.client.dsl.Resource;
import io.fabric8.kubernetes.client.server.mock.KubernetesServer;
import io.fabric8.kubernetes.internal.KubernetesDeserializer;
import io.fabric8.kubernetes.model.annotation.Group;
import io.fabric8.kubernetes.model.annotation.Kind;
import io.fabric8.kubernetes.model.annotation.Version;
import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.DisplayName;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.RegisterExtension;
import java.net.HttpURLConnection;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;

public class TestKafkaHelperUtilFourth {

  // JUnit5使用@RegisterExtension替代@Rule启动Mock Server
  @RegisterExtension
  public KubernetesServer server = new KubernetesServer(true, true);

  @Test
  @DisplayName("Should watch all custom resources")
  public void testWatch() throws InterruptedException {
    // 提前注册自定义资源类型,保证Mock序列化/反序列化正常
    KubernetesDeserializer.registerCustomKind("custom.example.com/v1", "UserACL", UserACL.class);
    
    // 用通配符匹配Watch路径,忽略自动拼接的resourceVersion等参数
    server.expect().withPathMatching("/apis/custom.example.com/v1/namespaces/default/useracls\\?watch.*")
        .andUpgradeToWebSocket()
        .open()
        .waitFor(100L)
        .andEmit(new WatchEvent(getUserACL("test-resource"), "ADDED"))
        .waitFor(100L)
        .andEmit(new WatchEventBuilder()
            .withNewStatusObject()
            .withMessage("410 - the event requested is outdated")
            .withCode(HttpURLConnection.HTTP_GONE)
            .endStatusObject()
            .build()).done().always();
    KubernetesClient client = server.getClient();
    MixedOperation<
        UserACL,
        KubernetesResourceList<UserACL>,
        Resource<UserACL>>
        userAclClient = client.resources(UserACL.class);

    CountDownLatch eventReceived = new CountDownLatch(1);
    Watch watch = userAclClient.inNamespace("default").watch(new Watcher<UserACL>() {
      @Override
      public void eventReceived(Action action, UserACL userAcl) {
        if (action == Action.ADDED) {
          eventReceived.countDown();
        }
      }

      @Override
      public void onClose(WatcherException e) { }
    });

    boolean eventArrived = eventReceived.await(5, TimeUnit.SECONDS);
    Assertions.assertTrue(eventArrived);
    watch.close();
  }

  private UserACL getUserACL(String resourceName) {
    UserACLSpec spec = new UserACLSpec();
    spec.setUserName("test-user-name");

    UserACL createdUserACL = new UserACL();
    createdUserACL.setMetadata(
        new ObjectMetaBuilder().withName(resourceName).build());
    createdUserACL.setSpec(spec);

    Condition condition = new Condition();
    condition.setMessage("Last reconciliation succeeded");
    condition.setReason("Successful");
    condition.setStatus("True");
    condition.setType("Successful");
    UserACLStatus status = new UserACLStatus();
    status.setCondition(new Condition[]{condition});
    createdUserACL.setStatus(status);

    return createdUserACL;
  }

  @Group("custom.example.com")
  @Version("v1")
  @Kind("UserACL")
  public static final class UserACL
      extends CustomResource<UserACLSpec, UserACLStatus> {
  }

  public static final class UserACLSpec {
    private String userName;

    public UserACLSpec() {}

    public String getUserName() {
      return userName;
    }

    public void setUserName(String userName) {
      this.userName = userName;
    }
  }

  public static final class UserACLStatus {
    Condition[] condition;

    public UserACLStatus() {};

    public Condition[] getCondition() {
      return condition;
    }

    public void setCondition(Condition[] condition) {
      this.condition = condition;
    }
  }
}

额外注意点

  • 构建配置中直接删除implementation group: 'io.fabric8', name: 'kubernetes-api', version: '3.0.12'依赖,避免类冲突。
  • 原代码waitFor(10L)仅等待10毫秒,在IO负载高的环境下可能出现时序问题,调整为100毫秒更稳妥。
  • 判断事件类型直接用枚举值比较action == Action.ADDED比字符串包含更严谨,避免误判。

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

相关产品推荐
方舟 Agent Plan

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

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