diff --git a/dubbo-dependencies-bom/pom.xml b/dubbo-dependencies-bom/pom.xml
index 8cc1d43d3d..8891d21bfd 100644
--- a/dubbo-dependencies-bom/pom.xml
+++ b/dubbo-dependencies-bom/pom.xml
@@ -148,7 +148,7 @@
1.9.12
- 4.10.3
+ 5.3.0
1.0.8
diff --git a/dubbo-registry/dubbo-registry-kubernetes/src/main/java/org/apache/dubbo/registry/kubernetes/KubernetesServiceDiscovery.java b/dubbo-registry/dubbo-registry-kubernetes/src/main/java/org/apache/dubbo/registry/kubernetes/KubernetesServiceDiscovery.java
index f75837ae4e..9719820baa 100644
--- a/dubbo-registry/dubbo-registry-kubernetes/src/main/java/org/apache/dubbo/registry/kubernetes/KubernetesServiceDiscovery.java
+++ b/dubbo-registry/dubbo-registry-kubernetes/src/main/java/org/apache/dubbo/registry/kubernetes/KubernetesServiceDiscovery.java
@@ -34,13 +34,14 @@ import io.fabric8.kubernetes.api.model.EndpointPort;
import io.fabric8.kubernetes.api.model.EndpointSubset;
import io.fabric8.kubernetes.api.model.Endpoints;
import io.fabric8.kubernetes.api.model.Pod;
+import io.fabric8.kubernetes.api.model.PodBuilder;
import io.fabric8.kubernetes.api.model.Service;
import io.fabric8.kubernetes.client.Config;
import io.fabric8.kubernetes.client.DefaultKubernetesClient;
import io.fabric8.kubernetes.client.KubernetesClient;
-import io.fabric8.kubernetes.client.KubernetesClientException;
import io.fabric8.kubernetes.client.Watch;
import io.fabric8.kubernetes.client.Watcher;
+import io.fabric8.kubernetes.client.WatcherException;
import java.util.HashSet;
import java.util.LinkedList;
@@ -123,11 +124,12 @@ public class KubernetesServiceDiscovery implements ServiceDiscovery {
.pods()
.inNamespace(namespace)
.withName(currentHostname)
- .edit()
- .editOrNewMetadata()
- .addToAnnotations(KUBERNETES_PROPERTIES_KEY, JSONObject.toJSONString(serviceInstance.getMetadata()))
- .endMetadata()
- .done();
+ .edit(pod->
+ new PodBuilder(pod)
+ .editOrNewMetadata()
+ .addToAnnotations(KUBERNETES_PROPERTIES_KEY, JSONObject.toJSONString(serviceInstance.getMetadata()))
+ .endMetadata()
+ .build());
if (logger.isInfoEnabled()) {
logger.info("Write Current Service Instance Metadata to Kubernetes pod. " +
"Current pod name: " + currentHostname);
@@ -149,11 +151,12 @@ public class KubernetesServiceDiscovery implements ServiceDiscovery {
.pods()
.inNamespace(namespace)
.withName(currentHostname)
- .edit()
- .editOrNewMetadata()
- .removeFromAnnotations(KUBERNETES_PROPERTIES_KEY)
- .endMetadata()
- .done();
+ .edit(pod ->
+ new PodBuilder(pod)
+ .editOrNewMetadata()
+ .removeFromAnnotations(KUBERNETES_PROPERTIES_KEY)
+ .endMetadata()
+ .build());
if (logger.isInfoEnabled()) {
logger.info("Remove Current Service Instance from Kubernetes pod. Current pod name: " + currentHostname);
}
@@ -222,7 +225,7 @@ public class KubernetesServiceDiscovery implements ServiceDiscovery {
}
@Override
- public void onClose(KubernetesClientException cause) {
+ public void onClose(WatcherException cause) {
// ignore
}
});
@@ -253,7 +256,7 @@ public class KubernetesServiceDiscovery implements ServiceDiscovery {
}
@Override
- public void onClose(KubernetesClientException cause) {
+ public void onClose(WatcherException cause) {
// ignore
}
});
@@ -284,7 +287,7 @@ public class KubernetesServiceDiscovery implements ServiceDiscovery {
}
@Override
- public void onClose(KubernetesClientException cause) {
+ public void onClose(WatcherException cause) {
// ignore
}
});
diff --git a/dubbo-registry/dubbo-registry-kubernetes/src/test/java/org/apache/dubbo/registry/kubernetes/KubernetesServiceDiscoveryTest.java b/dubbo-registry/dubbo-registry-kubernetes/src/test/java/org/apache/dubbo/registry/kubernetes/KubernetesServiceDiscoveryTest.java
index d3d8c9bd6e..15e9b1a88e 100644
--- a/dubbo-registry/dubbo-registry-kubernetes/src/test/java/org/apache/dubbo/registry/kubernetes/KubernetesServiceDiscoveryTest.java
+++ b/dubbo-registry/dubbo-registry-kubernetes/src/test/java/org/apache/dubbo/registry/kubernetes/KubernetesServiceDiscoveryTest.java
@@ -30,7 +30,7 @@ import io.fabric8.kubernetes.api.model.PodBuilder;
import io.fabric8.kubernetes.api.model.Service;
import io.fabric8.kubernetes.api.model.ServiceBuilder;
import io.fabric8.kubernetes.client.Config;
-import io.fabric8.kubernetes.client.KubernetesClient;
+import io.fabric8.kubernetes.client.NamespacedKubernetesClient;
import io.fabric8.kubernetes.client.server.mock.KubernetesServer;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.Assertions;
@@ -47,9 +47,9 @@ import java.util.Map;
@ExtendWith({MockitoExtension.class})
public class KubernetesServiceDiscoveryTest {
- public KubernetesServer mockServer = new KubernetesServer(true, true);
+ public KubernetesServer mockServer = new KubernetesServer(false, true);
- private KubernetesClient mockClient;
+ private NamespacedKubernetesClient mockClient;
private ServiceInstancesChangedListener mockListener = Mockito.mock(ServiceInstancesChangedListener.class);
@@ -64,7 +64,7 @@ public class KubernetesServiceDiscoveryTest {
serverUrl = URL.valueOf(mockClient.getConfiguration().getMasterUrl())
.setProtocol("kubernetes")
- .addParameter(KubernetesClientConst.TRUST_CERTS, "true")
+ .addParameter(KubernetesClientConst.USE_HTTPS, "false")
.addParameter(KubernetesClientConst.HTTP2_DISABLE, "true");
System.setProperty(Config.KUBERNETES_AUTH_TRYKUBECONFIG_SYSTEM_PROPERTY, "false");
@@ -116,16 +116,20 @@ public class KubernetesServiceDiscoveryTest {
Mockito.doNothing().when(mockListener).onEvent(Mockito.any());
serviceDiscovery.addServiceInstancesChangedListener(mockListener);
- mockClient.endpoints().withName("TestService").edit().editFirstSubset()
- .addNewAddress().withIp("ip2")
- .withNewTargetRef().withUid("uid2").withName("TestServer").endTargetRef().endAddress()
- .addNewPort("Test", "Test", 12345, "TCP").endSubset()
- .done();
+ mockClient.endpoints().withName("TestService")
+ .edit(endpoints ->
+ new EndpointsBuilder(endpoints)
+ .editFirstSubset()
+ .addNewAddress()
+ .withIp("ip2")
+ .withNewTargetRef().withUid("uid2").withName("TestServer").endTargetRef()
+ .endAddress().endSubset()
+ .build());
Thread.sleep(5000);
ArgumentCaptor eventArgumentCaptor =
ArgumentCaptor.forClass(ServiceInstancesChangedEvent.class);
- Mockito.verify(mockListener).onEvent(eventArgumentCaptor.capture());
+ Mockito.verify(mockListener, Mockito.times(2)).onEvent(eventArgumentCaptor.capture());
Assertions.assertEquals(2, eventArgumentCaptor.getValue().getServiceInstances().size());
serviceDiscovery.unregister(serviceInstance);
@@ -158,7 +162,7 @@ public class KubernetesServiceDiscoveryTest {
Thread.sleep(5000);
ArgumentCaptor eventArgumentCaptor =
ArgumentCaptor.forClass(ServiceInstancesChangedEvent.class);
- Mockito.verify(mockListener).onEvent(eventArgumentCaptor.capture());
+ Mockito.verify(mockListener, Mockito.times(2)).onEvent(eventArgumentCaptor.capture());
Assertions.assertEquals(1, eventArgumentCaptor.getValue().getServiceInstances().size());
serviceDiscovery.unregister(serviceInstance);