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
相关产品推荐
相关产品推荐

