feat(master): Optimize logic of uid calculate for better performance

1.use channel to instead of lock and isolate with main heartbeat routine
2.calculate in an aysnchoronus way periodically instead of real time
3.move heartbeat outside the goroutine, the startup of goroutine may delay the response time

Signed-off-by: leonrayang <changliang@oppo.com>
This commit is contained in:
leonrayang 2024-03-27 23:46:47 +08:00 committed by longerfly
parent 2b26a1fbcb
commit 4ee2208263
7 changed files with 132 additions and 48 deletions

View File

@ -126,7 +126,7 @@ func newUidDelCmd(client *master.MasterClient) *cobra.Command {
errout(err)
}()
var uidInfo *proto.UidSpaceRsp
if uidInfo, err = client.UserAPI().UidOperation(args[0], args[1], util.UidDel, ""); err != nil || !uidInfo.OK {
if uidInfo, err = client.UserAPI().UidOperation(args[0], args[1], util.UidDelLimit, ""); err != nil || !uidInfo.OK {
return
}
stdout("success!\n")

View File

@ -672,11 +672,14 @@ func (m *Server) UidOperate(w http.ResponseWriter, r *http.Request) {
case util.UidGetLimit:
ok, uidInfo = vol.uidSpaceManager.checkUid(uid)
uidList = append(uidList, uidInfo)
case util.AclAddIP:
ok = vol.uidSpaceManager.addUid(uid, capSize)
case util.AclDelIP:
ok = vol.uidSpaceManager.removeUid(uid)
case util.AclListIP:
case util.UidAddLimit, util.UidDelLimit:
cmd := &UidCmd{
op: op,
uid: uid,
size: capSize,
}
ok = vol.uidSpaceManager.pushUidCmd(cmd)
case util.UidLimitList:
uidList = vol.uidSpaceManager.listAll()
default:
// do nothing
@ -2680,8 +2683,8 @@ func newSimpleView(vol *Vol) (view *proto.SimpleVolView) {
DeleteExecTime: vol.DeleteExecTime,
}
vol.uidSpaceManager.RLock()
defer vol.uidSpaceManager.RUnlock()
vol.uidSpaceManager.rwMutex.RLock()
defer vol.uidSpaceManager.rwMutex.RUnlock()
for _, uid := range vol.uidSpaceManager.uidInfo {
view.Uids = append(view.Uids, proto.UidSimpleInfo{
UID: uid.Uid,

View File

@ -387,6 +387,7 @@ func (c *Cluster) scheduleTask() {
c.scheduleToLcScan()
c.scheduleToSnapshotDelVerScan()
c.scheduleToBadDisk()
c.scheduleToCheckVolUid()
}
func (c *Cluster) masterAddr() (addr string) {
@ -572,6 +573,23 @@ func (c *Cluster) scheduleToCheckVolQos() {
}()
}
func (c *Cluster) scheduleToCheckVolUid() {
go func() {
//check vols after switching leader two minutes
for {
if c.partition.IsRaftLeader() {
vols := c.copyVols()
for _, vol := range vols {
vol.uidSpaceManager.scheduleUidUpdate()
vol.uidSpaceManager.reCalculate()
}
}
// time.Sleep(time.Second * time.Duration(c.cfg.IntervalToCheckQos))
time.Sleep(time.Duration(float32(time.Second) * 0.5))
}
}()
}
func (c *Cluster) scheduleToCheckNodeSetGrpManagerStatus() {
go func() {
for {

View File

@ -1132,7 +1132,7 @@ func (c *Cluster) updateMetaNode(metaNode *MetaNode, metaPartitions []*proto.Met
}
mp.updateMetaPartition(mr, metaNode)
vol.uidSpaceManager.volUidUpdate(mr)
vol.uidSpaceManager.pushUidMsg(mr)
vol.quotaManager.quotaUpdate(mr)
c.updateInodeIDUpperBound(mp, mr, threshold, metaNode)
}

View File

@ -17,13 +17,24 @@ type UidSpaceManager struct {
uidInfo map[uint32]*proto.UidSpaceInfo
c *Cluster
vol *Vol
sync.RWMutex
msgChan chan *proto.MetaPartitionReport
cmdChan chan *UidCmd
exitC chan struct{}
rwMutex sync.RWMutex
}
type UidSpaceFsm struct {
UidSpaceArr []*proto.UidSpaceInfo
}
type UidCmd struct {
op uint64
uid uint32
size uint64
wg sync.WaitGroup
rsp interface{}
}
func (vol *Vol) initUidSpaceManager(c *Cluster) {
vol.uidSpaceManager = &UidSpaceManager{
c: c,
@ -31,49 +42,66 @@ func (vol *Vol) initUidSpaceManager(c *Cluster) {
volName: vol.Name,
mpSpaceMetrics: make(map[uint64][]*proto.UidReportSpaceInfo),
uidInfo: make(map[uint32]*proto.UidSpaceInfo),
msgChan: make(chan *proto.MetaPartitionReport, 10000),
cmdChan: make(chan *UidCmd, 1000),
}
}
func (uMgr *UidSpaceManager) addUid(uid uint32, size uint64) bool {
uMgr.Lock()
uMgr.uidInfo[uid] = &proto.UidSpaceInfo{
LimitSize: size,
VolName: uMgr.volName,
Uid: uid,
Enabled: true,
func (uMgr *UidSpaceManager) pushUidCmd(cmd *UidCmd) bool {
cmd.wg.Add(1)
select {
case uMgr.cmdChan <- cmd:
default:
log.LogWarnf("vol %v volUidUpdate.mpID %v uid %v op %v be missed", uMgr.volName, cmd.uid, cmd.op)
return false
}
uMgr.persist()
uMgr.Unlock()
uMgr.listAll()
log.LogDebugf("pushUidCmd. vol %v cmd (%v) wait result", uMgr.volName, cmd)
cmd.wg.Wait()
log.LogDebugf("pushUidCmd. vol %v cmd (%v) get result", uMgr.volName, cmd)
return true
}
func (uMgr *UidSpaceManager) removeUid(uid uint32) bool {
uMgr.Lock()
defer uMgr.Unlock()
func (uMgr *UidSpaceManager) addUid(cmd *UidCmd) {
defer cmd.wg.Done()
if uidInfo, ok := uMgr.uidInfo[cmd.uid]; ok {
if uidInfo.Enabled == true {
log.LogWarnf("UidSpaceManager.addUid vol %v add %v already exist", uMgr.volName, cmd.uid)
return
}
}
uMgr.uidInfo[cmd.uid] = &proto.UidSpaceInfo{
LimitSize: cmd.size,
VolName: uMgr.volName,
Uid: cmd.uid,
Enabled: true,
}
uMgr.persist()
log.LogWarnf("UidSpaceManager.vol %v addUid %v success", uMgr.volName, cmd.uid)
}
if _, ok := uMgr.uidInfo[uid]; !ok {
log.LogErrorf("UidSpaceManager.vol %v del %v failed", uMgr.volName, uid)
func (uMgr *UidSpaceManager) removeUid(cmd *UidCmd) bool {
defer cmd.wg.Done()
if _, ok := uMgr.uidInfo[cmd.uid]; !ok {
log.LogWarnf("UidSpaceManager.vol %v uid del uid %v not exist", uMgr.volName, cmd.uid)
return true
}
uMgr.uidInfo[uid].Enabled = false
uMgr.uidInfo[uid].Limited = false
uMgr.uidInfo[cmd.uid].Enabled = false
uMgr.uidInfo[cmd.uid].Limited = false
uMgr.persist()
log.LogDebugf("UidSpaceManager.vol %v del %v success", uMgr.volName, uid)
log.LogWarnf("UidSpaceManager.vol %v del %v success", uMgr.volName, cmd.uid)
return true
}
func (uMgr *UidSpaceManager) checkUid(uid uint32) (ok bool, uidInfo *proto.UidSpaceInfo) {
uMgr.RLock()
defer uMgr.RUnlock()
uMgr.rwMutex.RLock()
defer uMgr.rwMutex.RUnlock()
uidInfo, ok = uMgr.uidInfo[uid]
return
}
func (uMgr *UidSpaceManager) listAll() (rsp []*proto.UidSpaceInfo) {
uMgr.RLock()
defer uMgr.RUnlock()
uMgr.rwMutex.RLock()
defer uMgr.rwMutex.RUnlock()
log.LogDebugf("UidSpaceManager. listAll vol %v, info %v", uMgr.volName, len(uMgr.uidInfo))
for _, t := range uMgr.uidInfo {
@ -118,8 +146,8 @@ func (uMgr *UidSpaceManager) load(c *Cluster, val []byte) (err error) {
}
func (uMgr *UidSpaceManager) getSpaceOp() (rsp []*proto.UidSpaceInfo) {
uMgr.RLock()
defer uMgr.RUnlock()
uMgr.rwMutex.RLock()
defer uMgr.rwMutex.RUnlock()
for _, info := range uMgr.uidInfo {
rsp = append(rsp, info)
log.LogDebugf("getSpaceOp. vol %v uid %v enabled %v", info.VolName, info.Uid, info.Limited)
@ -127,15 +155,43 @@ func (uMgr *UidSpaceManager) getSpaceOp() (rsp []*proto.UidSpaceInfo) {
return
}
func (uMgr *UidSpaceManager) volUidUpdate(report *proto.MetaPartitionReport) {
func (uMgr *UidSpaceManager) scheduleUidUpdate() {
for {
select {
case report := <-uMgr.msgChan:
log.LogDebugf("vol %v scheduleUidUpdate.mpID %v set uid %v be set", uMgr.volName, report.PartitionID, report.UidInfo)
uMgr.volUidUpdate(report)
case cmd := <-uMgr.cmdChan:
uMgr.rwMutex.Lock()
log.LogDebugf("vol %v scheduleUidUpdate.cmd(%v)", uMgr.volName, cmd)
if cmd.op == util.UidAddLimit {
uMgr.addUid(cmd)
} else if cmd.op == util.UidDelLimit {
uMgr.removeUid(cmd)
}
uMgr.rwMutex.Unlock()
log.LogDebugf("vol %v scheduleUidUpdate.cmd(%v) left", uMgr.volName, cmd)
default:
return
}
}
}
func (uMgr *UidSpaceManager) pushUidMsg(report *proto.MetaPartitionReport) {
if !report.IsLeader {
return
}
uMgr.Lock()
defer uMgr.Unlock()
id := report.PartitionID
uMgr.mpSpaceMetrics[id] = report.UidInfo
log.LogDebugf("vol %v volUidUpdate.mpID %v set uid %v. uid list size %v", uMgr.volName, id, report.UidInfo, len(uMgr.uidInfo))
select {
case uMgr.msgChan <- report:
default:
log.LogWarnf("vol %v volUidUpdate.mpID %v set uid %v be missed", uMgr.volName, report.PartitionID, report.UidInfo)
return
}
}
func (uMgr *UidSpaceManager) reCalculate() {
uMgr.rwMutex.Lock()
defer uMgr.rwMutex.Unlock()
for _, info := range uMgr.uidInfo {
info.UsedSize = 0
@ -143,7 +199,6 @@ func (uMgr *UidSpaceManager) volUidUpdate(report *proto.MetaPartitionReport) {
uidInfo := make(map[uint32]*proto.UidSpaceInfo)
for mpId, info := range uMgr.mpSpaceMetrics {
log.LogDebugf("vol %v volUidUpdate. reCalc mpId %v info %v", uMgr.volName, mpId, len(info))
for _, space := range info {
if _, ok := uMgr.uidInfo[space.Uid]; !ok {
log.LogDebugf("vol %v volUidUpdate.uid %v not found", uMgr.volName, space.Uid)
@ -158,7 +213,7 @@ func (uMgr *UidSpaceManager) volUidUpdate(report *proto.MetaPartitionReport) {
uidInfo[space.Uid] = &infoCopy
}
log.LogDebugf("volUidUpdate.vol %v uid %v from mpId %v useSize %v add %v", uMgr.vol, space.Uid, mpId, uidInfo[space.Uid].UsedSize, space.Size)
// log.LogDebugf("volUidUpdate.vol %v uid %v from mpId %v useSize %v add %v", uMgr.vol, space.Uid, mpId, uidInfo[space.Uid].UsedSize, space.Size)
uidInfo[space.Uid].UsedSize += space.Size
if !uidInfo[space.Uid].Enabled {
uidInfo[space.Uid].Limited = false
@ -174,7 +229,6 @@ func (uMgr *UidSpaceManager) volUidUpdate(report *proto.MetaPartitionReport) {
}
}
log.LogDebugf("vol %v volUidUpdate.mpID %v set uid %v. uid list size %v", uMgr.volName, id, report.UidInfo, len(uMgr.uidInfo))
for _, info := range uidInfo {
if _, ok := uMgr.uidInfo[info.Uid]; !ok {
log.LogErrorf("volUidUpdate.uid %v not found", info.Uid)
@ -187,7 +241,17 @@ func (uMgr *UidSpaceManager) volUidUpdate(report *proto.MetaPartitionReport) {
info.Limited = false
}
}
log.LogDebugf("volUidUpdate.mpID %v set uid %v. uid list size %v", id, report.UidInfo, len(uMgr.uidInfo))
// log.LogDebugf("volUidUpdate. uid count %v", len(uMgr.uidInfo))
}
func (uMgr *UidSpaceManager) volUidUpdate(report *proto.MetaPartitionReport) {
if !report.IsLeader {
return
}
id := report.PartitionID
uMgr.mpSpaceMetrics[id] = report.UidInfo
log.LogDebugf("vol %v volUidUpdate.mpID %v set uid %v. uid list size %v", uMgr.volName, id, report.UidInfo, len(uMgr.uidInfo))
}
type ServerFactorLimit struct {

View File

@ -83,9 +83,8 @@ func (m *metadataManager) opMasterHeartbeat(conn net.Conn, p *Packet, remoteAddr
Request: req,
}
)
start := time.Now()
go func() {
start := time.Now()
decode := json.NewDecoder(bytes.NewBuffer(data))
decode.UseNumber()
if err = decode.Decode(adminTask); err != nil {

View File

@ -65,7 +65,7 @@ const (
const (
UidLimitList = 0
UidAddLimit = 1
UidDel = 2
UidDelLimit = 2
UidGetLimit = 3
)