diff --git a/dubbo-xds/src/main/java/org/apache/dubbo/registry/xds/util/protocol/AbstractProtocol.java b/dubbo-xds/src/main/java/org/apache/dubbo/registry/xds/util/protocol/AbstractProtocol.java index 9b6687f54b..eb21cf3548 100644 --- a/dubbo-xds/src/main/java/org/apache/dubbo/registry/xds/util/protocol/AbstractProtocol.java +++ b/dubbo-xds/src/main/java/org/apache/dubbo/registry/xds/util/protocol/AbstractProtocol.java @@ -154,6 +154,11 @@ public abstract class AbstractProtocol> implements // observer reused StreamObserver observer = requestObserverMap.get(request); + if (observer == null) { + observer = xdsChannel.createDeltaDiscoveryRequest(new ResponseObserver(request)); + requestObserverMap.put(request, observer); + } + // send request to control panel observer.onNext(buildDiscoveryRequest(names)); @@ -231,11 +236,17 @@ public abstract class AbstractProtocol> implements @Override public void onError(Throwable t) { logger.error("xDS Client received error message! detail:", t); + clear(); } @Override public void onCompleted() { - // ignore + logger.info("xDS Client completed, requestId: " + requestId); + clear(); + } + + private void clear() { + requestObserverMap.remove(requestId); } }