etcd: fix "fatal error: concurrent map writes"

Trace

	goroutine 1270434 [running]:
	runtime.throw(0x4b0cc5e, 0x15)
		/home/yousong/.usr/go/goroot-1.15.6/src/runtime/panic.go:1116 +0x72 fp=0xc0011151f8 sp=0xc0011151c8 pc=0x439292
	runtime.mapdelete_faststr(0x40876e0, 0xc0006a54a0, 0xc001e84200, 0x3e)
		/home/yousong/.usr/go/goroot-1.15.6/src/runtime/map_faststr.go:377 +0x34c fp=0xc001115260 sp=0xc0011151f8 pc=0x41664c
	yunion.io/x/onecloud/pkg/cloudcommon/etcd.(*SEtcdClient).Unwatch(0xc0003130e0, 0xc001e84200, 0x3e)
		/home/yousong/go/src/yunion.io/x/onecloud/pkg/cloudcommon/etcd/etcd.go:353 +0x134 fp=0xc0011152d0 sp=0xc001115260 pc=0x10031f4
	yunion.io/x/onecloud/pkg/compute/models.(*SHostHealthChecker).WatchHost(0xc000c222a0, 0x5229ca0, 0xc000122000, 0xc001fbc4e0, 0x24)
		/home/yousong/go/src/yunion.io/x/onecloud/pkg/compute/models/host_health.go:149 +0xd5 fp=0xc001115338 sp=0xc0011152d0 pc=0x1816395
	yunion.io/x/onecloud/pkg/compute/models.(*SHost).PerformOnline(0xc001fd8400, 0x5229d20, 0xc001c4d440, 0x52a0c60, 0xc000957c40, 0x5298900, 0xc001629aa0, 0x5298900, 0xc001629ae0, 0xc000957c40, ...)
		/home/yousong/go/src/yunion.io/x/onecloud/pkg/compute/models/hosts.go:3829 +0x227 fp=0xc0011153a0 sp=0xc001115338 pc=0x185a387
	yunion.io/x/onecloud/pkg/compute/models.(*SHost).PerformPing(0xc001fd8400, 0x5229d20, 0xc001c4d440, 0x52a0c60, 0xc000957c40, 0x5298900, 0xc001629aa0, 0x5298900, 0xc001629ae0, 0x0, ...)
		/home/yousong/go/src/yunion.io/x/onecloud/pkg/compute/models/hosts.go:3891 +0xbd fp=0xc001115430 sp=0xc0011153a0 pc=0x185acbd
This commit is contained in:
Yousong Zhou
2021-03-08 09:47:50 +08:00
parent cdde5b31de
commit b424a8f67d
+9 -1
View File
@@ -18,6 +18,7 @@ import (
"context"
"crypto/tls"
"fmt"
"sync"
"time"
"go.etcd.io/etcd/clientv3"
@@ -45,7 +46,8 @@ type SEtcdClient struct {
onKeepaliveFailure func()
leaseLiving bool
watchers map[string]*SEtcdWatcher
watchers map[string]*SEtcdWatcher
watchersMu *sync.Mutex
}
func defaultOnKeepAliveFailed() {
@@ -115,6 +117,7 @@ func NewEtcdClient(opt *SEtcdOptions, onKeepaliveFailure func()) (*SEtcdClient,
etcdClient.leaseTtlTimeout = timeoutSeconds
etcdClient.watchers = make(map[string]*SEtcdWatcher)
etcdClient.watchersMu = &sync.Mutex{}
etcdClient.namespace = opt.EtcdNamspace
@@ -306,8 +309,10 @@ func (cli *SEtcdClient) Watch(
onModify TEtcdModifyEventFunc,
onDelete TEtcdDeleteEventFunc,
) error {
cli.watchersMu.Lock()
_, ok := cli.watchers[prefix]
if ok {
cli.watchersMu.Unlock()
return errors.Errorf("watch prefix %s already registered", prefix)
}
@@ -318,6 +323,7 @@ func (cli *SEtcdClient) Watch(
watcher: watcher,
cancel: cancel,
}
cli.watchersMu.Unlock()
prefix = cli.getKey(prefix)
@@ -346,6 +352,8 @@ func (cli *SEtcdClient) Watch(
}
func (cli *SEtcdClient) Unwatch(prefix string) {
cli.watchersMu.Lock()
defer cli.watchersMu.Unlock()
watcher, ok := cli.watchers[prefix]
if ok {
log.Debugf("unwatch %s", prefix)