diff options
| author | KubeEdge Bot <48982446+kubeedge-bot@users.noreply.github.com> | 2022-03-15 13:42:56 +0800 |
|---|---|---|
| committer | GitHub <noreply@github.com> | 2022-03-15 13:42:56 +0800 |
| commit | effcef246ea18dda3524ab3d16c4a2dda56deae0 (patch) | |
| tree | 25826c5885e3301db1f111fbadda64116e49da1d | |
| parent | Merge pull request #3703 from Shelley-BaoYue/automated-cherry-pick-of-#3693-u... (diff) | |
| parent | fix concurrent map iteration and map write bug (diff) | |
| download | kubeedge-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.go | 16 |
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()) } |
