FoghostCn commented on code in PR #2859:
URL: https://github.com/apache/dubbo-go/pull/2859#discussion_r2081429932
##########
registry/nacos/registry.go:
##########
@@ -177,46 +177,67 @@ func (nr *nacosRegistry) Subscribe(url *common.URL,
notifyListener registry.Noti
return nil
}
serviceName := url.GetParam(constant.InterfaceKey, "")
- var serviceNames []string
- var err error
if serviceName == constant.AnyValue {
- serviceNames, err = nr.getAllSubscribeServiceNames(url)
- if err != nil {
- return err
- }
+ nr.subscribeAll(url, notifyListener)
+ go func() {
+ // scheduled lookup for new service
+ for {
+ nr.subscribeAll(url, notifyListener)
+ time.Sleep(LookupInterval)
+ }
+ }()
+ return nil
} else {
- serviceNames = []string{getSubscribeName(url)}
+ // retry forever
+ for {
+ err := nr.subscribe(getSubscribeName(url),
notifyListener)
+ if err == nil {
+ return nil
+ }
+ }
}
- return nr.subscribe(serviceNames, notifyListener)
}
-// subscribe subscribe services
-func (nr *nacosRegistry) subscribe(serviceNames []string, notifyListener
registry.NotifyListener) error {
+func (nr *nacosRegistry) subscribeAll(url *common.URL, notifyListener
registry.NotifyListener) {
+ groupName := nr.URL.GetParam(constant.RegistryGroupKey, defaultGroup)
+ serviceNames, err := nr.getAllSubscribeServiceNames(url)
+ if err != nil {
+ logger.Warnf("getAllServices() = err:%v",
perrors.WithStack(err))
+ return
+ }
if len(serviceNames) == 0 {
logger.Warnf("No services to listen to.")
- return nil
+ return
}
- for {
- if !nr.IsAvailable() {
- logger.Warnf("event listener game over.")
- return perrors.New("nacosRegistry is not available.")
+ for _, name := range serviceNames {
+ if _, ok := listenerCache.Load(name + groupName); ok {
+ continue
}
- var err error
- for _, serviceName := range serviceNames {
- listener :=
NewNacosListenerWithServiceName(serviceName, nr.URL, nr.namingClient)
- err = listener.listenService(serviceName)
- metrics.Publish(metricsRegistry.NewSubscribeEvent(err
== nil))
- if err != nil {
- logger.Warnf("getAllServices() = err:%v",
perrors.WithStack(err))
- time.Sleep(time.Duration(RegistryConnDelay) *
time.Second)
- break
- }
- go nr.handleServiceEvents(listener, notifyListener)
- }
- if err == nil {
- break
+ err = nr.subscribe(name, notifyListener)
+ if err != nil {
+ logger.Warnf("subscribe service %s err:%v", name,
perrors.WithStack(err))
}
}
+}
+
+// subscribe subscribe services
+func (nr *nacosRegistry) subscribe(serviceName string, notifyListener
registry.NotifyListener) error {
+ if len(serviceName) == 0 {
+ logger.Warnf("can not subscribe because service name is empty")
+ return nil
+ }
+ if !nr.IsAvailable() {
+ logger.Warnf("event listener game over.")
+ return perrors.New("nacosRegistry is not available.")
+ }
+ listener := NewNacosListenerWithServiceName(serviceName, nr.URL,
nr.namingClient)
+ err := listener.listenService(serviceName)
+ metrics.Publish(metricsRegistry.NewSubscribeEvent(err == nil))
+ if err != nil {
+ logger.Warnf("subscribe service %s err:%v", serviceName,
perrors.WithStack(err))
+ return err
+ }
+ go nr.handleServiceEvents(listener, notifyListener)
Review Comment:
handleServiceEvents 中是有退出逻辑的,在这个registry的存续期间这个 gr 是需要一直存在的
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]