mirror of
https://github.com/cubefs/cubefs.git
synced 2026-08-02 02:00:56 +00:00
fix(flash): print hanging size
with: #1000307099
Signed-off-by: clinx <chenlin1@oppo.com>
(cherry picked from commit 7a249e1d99)
This commit is contained in:
parent
7f9177766d
commit
57bea86ae0
@ -77,6 +77,7 @@ const (
|
||||
_defaultMissEntryExpiration = 2 * time.Minute
|
||||
_defaultMaxMissEntryCache = 100000
|
||||
_defaultMissCountThresholdInterval = 5
|
||||
_defaultFlashLimitHangTimeout = 1000 // ms
|
||||
)
|
||||
|
||||
// Configuration keys
|
||||
@ -390,8 +391,8 @@ func (f *FlashNode) parseConfig(cfg *config.Config) (err error) {
|
||||
f.disks = disks
|
||||
}
|
||||
f.handleReadTimeout = proto.DefaultRemoteCacheHandleReadTimeout
|
||||
f.limitWrite = util.NewIOLimiterEx(f.diskWriteFlow, f.diskWriteIocc*len(f.disks), f.diskWriteIoFactorFlow, f.handleReadTimeout)
|
||||
f.limitRead = util.NewIOLimiterEx(f.diskReadFlow, f.diskReadIocc*len(f.disks), f.diskReadIoFactorFlow, f.handleReadTimeout)
|
||||
f.limitWrite = util.NewIOLimiterEx(f.diskWriteFlow, f.diskWriteIocc*len(f.disks), f.diskWriteIoFactorFlow, _defaultFlashLimitHangTimeout)
|
||||
f.limitRead = util.NewIOLimiterEx(f.diskReadFlow, f.diskReadIocc*len(f.disks), f.diskReadIoFactorFlow, _defaultFlashLimitHangTimeout)
|
||||
lruFhCapacity := cfg.GetInt(cfgLruFhCapacity)
|
||||
if lruFhCapacity <= 0 || lruFhCapacity >= 1000000 {
|
||||
lruFhCapacity = _defaultLRUFhCapacity
|
||||
|
||||
@ -288,6 +288,9 @@ func (f *FlashNode) opCacheDelete(conn net.Conn, p *proto.Packet) (err error) {
|
||||
return proto.ErrorNoCacheDeleteRequest
|
||||
}
|
||||
uniKey := string(data)
|
||||
if log.EnableDebug() {
|
||||
log.LogDebugf("action[opCacheDelete] delete key(%v) from remote(%v)", uniKey, conn.RemoteAddr().String())
|
||||
}
|
||||
pDir := cachengine.MapKeyToDirectory(uniKey)
|
||||
f.cacheEngine.DeleteCacheBlock(cachengine.GenCacheBlockKeyV2(pDir, uniKey))
|
||||
return nil
|
||||
|
||||
@ -649,7 +649,16 @@ func (rc *RemoteCacheClient) Put(ctx context.Context, reqId, key string, r io.Re
|
||||
return ch.RemoteError
|
||||
}
|
||||
currentReadSize := readSize - bufferOffset
|
||||
n, err = r.Read(buf[bufferOffset : bufferOffset+currentReadSize])
|
||||
readChan := make(chan struct{})
|
||||
go func() {
|
||||
n, err = r.Read(buf[bufferOffset : bufferOffset+currentReadSize])
|
||||
close(readChan)
|
||||
}()
|
||||
select {
|
||||
case <-readChan:
|
||||
case <-ctx.Done():
|
||||
return ctx.Err()
|
||||
}
|
||||
bufferOffset += n
|
||||
if err == io.EOF {
|
||||
break
|
||||
@ -1333,6 +1342,7 @@ func (rc *RemoteCacheClient) ReadObject(ctx context.Context, fg *FlashGroup, req
|
||||
}
|
||||
stat.EndStat("flashNode", err, bgTime, 1)
|
||||
}()
|
||||
var stcTime, etcTime, rcvTime int64
|
||||
for {
|
||||
if ctx.Err() != nil {
|
||||
return nil, 0, ctx.Err()
|
||||
@ -1349,6 +1359,7 @@ func (rc *RemoteCacheClient) ReadObject(ctx context.Context, fg *FlashGroup, req
|
||||
log.LogWarnf("%v FlashGroup Read: failed to MarshalData (%+v). err(%v)", logPrefix, req, err)
|
||||
return
|
||||
}
|
||||
stcTime = time.Now().UnixNano()
|
||||
if conn, err = rc.conns.GetConnect(addr); err != nil {
|
||||
log.LogWarnf("%v FlashGroup Read: get connection failed, addr(%v) reqPacket(%v) err(%v) remoteCacheMultiRead(%v)",
|
||||
logPrefix, addr, req, err, rc.RemoteCacheMultiRead)
|
||||
@ -1359,7 +1370,7 @@ func (rc *RemoteCacheClient) ReadObject(ctx context.Context, fg *FlashGroup, req
|
||||
}
|
||||
return
|
||||
}
|
||||
|
||||
etcTime = time.Now().UnixNano()
|
||||
if err = reqPacket.WriteToConn(conn); err != nil {
|
||||
log.LogWarnf("%v FlashGroup Read: failed to write to addr(%v) err(%v) remoteCacheMultiRead(%v)",
|
||||
logPrefix, addr, err, rc.RemoteCacheMultiRead)
|
||||
@ -1383,13 +1394,14 @@ func (rc *RemoteCacheClient) ReadObject(ctx context.Context, fg *FlashGroup, req
|
||||
fg.moveToUnknownRank(addr, err, rc.FlashNodeTimeoutCount)
|
||||
return nil, 0, err
|
||||
}
|
||||
rcvTime = time.Now().UnixNano()
|
||||
reader = rc.NewRemoteCacheReader(ctx, conn, reqId, flashIp, uint32(req.Offset), uint32(req.Offset+req.Size_))
|
||||
reader.needReadLen = int64(req.Size_)
|
||||
break
|
||||
}
|
||||
if log.EnableDebug() {
|
||||
log.LogDebugf("%v FlashGroup Read: flashGroup(%v) addr(%v) CacheReadRequest(%v) reqPacket(%v) err(%v)"+
|
||||
" remoteCacheMultiRead(%v)", logPrefix, fg, addr, req, reqPacket, err, rc.RemoteCacheMultiRead)
|
||||
" remoteCacheMultiRead(%v) getConn cost(%v) fisrtPackage cost(%v)", logPrefix, fg, addr, req, reqPacket, err, rc.RemoteCacheMultiRead, etcTime-stcTime, rcvTime-etcTime)
|
||||
}
|
||||
return
|
||||
}
|
||||
|
||||
@ -45,6 +45,7 @@ type LimiterStatus struct {
|
||||
IOQueue int
|
||||
IORunning int
|
||||
IOWaiting int
|
||||
IOHanging int
|
||||
Factor int
|
||||
}
|
||||
|
||||
@ -327,6 +328,7 @@ func (q *ioQueue) Status() (st LimiterStatus) {
|
||||
st.IOQueue = cap(q.queue)
|
||||
st.IORunning = int(atomic.LoadUint32(&q.running))
|
||||
st.IOWaiting = len(q.queue)
|
||||
st.IOHanging = len(q.midQueue)
|
||||
st.Factor = q.factor
|
||||
return
|
||||
}
|
||||
|
||||
Loading…
Reference in New Issue
Block a user