summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorKubeEdge Bot <48982446+kubeedge-bot@users.noreply.github.com>2022-03-15 13:42:56 +0800
committerGitHub <noreply@github.com>2022-03-15 13:42:56 +0800
commiteffcef246ea18dda3524ab3d16c4a2dda56deae0 (patch)
tree25826c5885e3301db1f111fbadda64116e49da1d
parentMerge pull request #3703 from Shelley-BaoYue/automated-cherry-pick-of-#3693-u... (diff)
parentfix concurrent map iteration and map write bug (diff)
downloadkubeedge-effcef246ea18dda3524ab3d16c4a2dda56deae0.tar.gz
Merge pull request #3710 from Shelley-BaoYue/automated-cherry-pick-of-#3670-upstream-release-1.10
Automated cherry pick of #3670: fix concurrent map iteration and map write bug
-rw-r--r--cloud/pkg/dynamiccontroller/application/eventhandler.go16
1 files changed, 12 insertions, 4 deletions
diff --git a/cloud/pkg/dynamiccontroller/application/eventhandler.go b/cloud/pkg/dynamiccontroller/application/eventhandler.go
index a4d754ce3..637f69767 100644
--- a/cloud/pkg/dynamiccontroller/application/eventhandler.go
+++ b/cloud/pkg/dynamiccontroller/application/eventhandler.go
@@ -18,6 +18,7 @@ package application
import (
"fmt"
+ "sync"
"k8s.io/apimachinery/pkg/runtime"
"k8s.io/apimachinery/pkg/runtime/schema"
@@ -82,10 +83,11 @@ func (c *handlerCenter) DeleteListener(s *SelectorListener) {
type CommonResourceEventHandler struct {
events chan watch.Event
//TODO: num of listeners is proportional to the number of request, need reduce.
- listeners map[string]*SelectorListener
- messageLayer messagelayer.MessageLayer
- gvr schema.GroupVersionResource
- informer informers.GenericInformer
+ listeners map[string]*SelectorListener
+ listenersLock sync.RWMutex
+ messageLayer messagelayer.MessageLayer
+ gvr schema.GroupVersionResource
+ informer informers.GenericInformer
}
func NewCommonResourceEventHandler(gvr schema.GroupVersionResource, informerFactory dynamicinformer.DynamicSharedInformerFactory, layer messagelayer.MessageLayer) *CommonResourceEventHandler {
@@ -144,20 +146,26 @@ func (c *CommonResourceEventHandler) AddListener(s *SelectorListener) error {
return fmt.Errorf("failed to list: %v", err)
}
s.sendAllObjects(ret, c.messageLayer)
+ c.listenersLock.Lock()
c.listeners[s.id] = s
+ c.listenersLock.Unlock()
return nil
}
func (c *CommonResourceEventHandler) DeleteListener(s *SelectorListener) {
+ c.listenersLock.Lock()
delete(c.listeners, s.id)
+ c.listenersLock.Unlock()
}
func (c *CommonResourceEventHandler) dispatchEvents() {
for event := range c.events {
klog.V(4).Infof("[metaserver/resourceEventHandler] handler(%v), send obj event{%v/%v} to listeners", c.gvr, event.Type, event.Object.GetObjectKind().GroupVersionKind().String())
+ c.listenersLock.RLock()
for _, listener := range c.listeners {
listener.sendObj(event, c.messageLayer)
}
+ c.listenersLock.RUnlock()
}
klog.Warningf("[metaserver/resourceEventHandler] handler(%v) stopped!", c.gvr.String())
}