fix(flash): warm up task set force remote cache true

with: #1000390069

Signed-off-by: clinx <chenlin1@oppo.com>
(cherry picked from commit 272660d7cb)
This commit is contained in:
clinx 2025-10-20 20:13:39 +08:00 committed by chihe
parent 9f3548653a
commit cd086fcb46
7 changed files with 53 additions and 20 deletions

View File

@ -12,7 +12,7 @@ import (
)
const (
minRequiredTTLSeconds = 30 * 60 // 30 minute
minRequiredTTLSeconds = 10 * 60 // 10 minute
FlashTaskInteractiveTimeout = time.Minute * time.Duration(5)
ScanManualInterval = time.Minute * time.Duration(2)
CheckInteractiveInterval = time.Minute * time.Duration(2)
@ -50,7 +50,7 @@ func checkManualConfig(flt *proto.FlashManualTask, vol *Vol, fltMgr *flashManual
switch flt.Action {
case proto.FlashManualWarmupAction:
if vol.remoteCacheTTL < minRequiredTTLSeconds && vol.remoteCacheTTL != 0 {
return fmt.Errorf("cache ttl %v of vol[%v] is too short and warm up can not work", vol.remoteCacheTTL, vol.Name)
return fmt.Errorf("cache ttl [%v]s of vol[%v] less than [%v]s and warm up can not work", vol.remoteCacheTTL, vol.Name, minRequiredTTLSeconds)
}
if err = isDuplicated(flt, fltMgr); err != nil {
return err

View File

@ -278,7 +278,9 @@ func (m *Server) createFlashNodeManualTask(w http.ResponseWriter, r *http.Reques
sendErrReply(w, r, &proto.HTTPReply{Code: proto.ErrCodeParamError, Msg: err.Error()})
return
}
if req.ManualTaskStatistics == nil {
req.ManualTaskStatistics = &proto.ManualTaskStatistics{}
}
// Validate file size limits
if req.ManualTaskConfig.MinFileSizeLimit > req.ManualTaskConfig.MaxFileSizeLimit {
err = fmt.Errorf("MinFileSizeLimit(%d) cannot be greater than MaxFileSizeLimit(%d)",

View File

@ -580,6 +580,9 @@ func (flt *FlashManualTask) GetPathPrefix() string {
func (flt *FlashManualTask) SetResponse(taskRsp *FlashNodeManualTaskResponse) {
t := time.Now()
flt.UpdateTime = &t
if flt.ManualTaskStatistics == nil {
flt.ManualTaskStatistics = &ManualTaskStatistics{}
}
flt.ManualTaskStatistics.FlashNode = taskRsp.FlashNode
flt.ManualTaskStatistics.TotalFileScannedNum = taskRsp.TotalFileScannedNum
flt.ManualTaskStatistics.TotalFileCachedNum = taskRsp.TotalFileCachedNum

View File

@ -844,6 +844,7 @@ func (f *FlashNode) getValidViewInfo(req *proto.FlashNodeManualTaskRequest) (met
MetaWrapper: metaWrapper,
NeedRemoteCache: true,
HeartBeatPing: true,
ForceRemoteCache: true,
}
log.LogInfof("[NewS3Scanner] extentConfig: vol(%v) volStorageClass(%v) allowedStorageClass(%v), followerRead(%v)",
extentConfig.Volume, extentConfig.VolStorageClass, extentConfig.VolAllowedStorageClass, extentConfig.FollowerRead)

View File

@ -927,6 +927,9 @@ func (client *ExtentClient) Close() error {
_ = client.EvictStream(inode)
}
client.dataWrapper.Stop()
if client.RemoteCache.Started {
client.RemoteCache.Stop()
}
return nil
}

View File

@ -117,6 +117,15 @@ type RemoteCacheClient struct {
WriteChunkSize int64 // bytes
}
func (rc *RemoteCacheClient) EnqueueConnTask(task *ConnPutTask) {
select {
case <-rc.stopC:
rc.processConnPut(task)
default:
rc.connPutChan <- task
}
}
func (as *AddressPingStats) Add(duration time.Duration) {
as.Lock()
defer as.Unlock()
@ -268,7 +277,22 @@ func (rc *RemoteCacheClient) consumeConnPut() {
for {
select {
case <-rc.stopC:
log.LogDebugf("consumeConnPut: stopping due to stopC signal")
log.LogDebugf("consumeConnPut: stopC received, draining connPutChan for up to 10s")
func() {
timer := time.NewTimer(10 * time.Second)
defer timer.Stop()
for {
select {
case connTask := <-rc.connPutChan:
if connTask != nil {
rc.processConnPut(connTask)
}
case <-timer.C:
log.LogDebugf("consumeConnPut: drain window elapsed, exiting")
return
}
}
}()
return
case connTask := <-rc.connPutChan:
if connTask != nil {
@ -584,7 +608,7 @@ func (rc *RemoteCacheClient) HeartBeat(addr string) (duration time.Duration, err
packet.Opcode = proto.OpFlashSDKHeartbeat
defer func() {
rc.connPutChan <- &ConnPutTask{conn: conn, forceClose: err != nil}
rc.EnqueueConnTask(&ConnPutTask{conn: conn, forceClose: err != nil})
}()
if conn, err = rc.conns.GetConnect(addr); err != nil {
@ -750,7 +774,7 @@ func (rc *RemoteCacheClient) Put(ctx context.Context, reqId, key string, r io.Re
bgTime := stat.BeginStat()
forceClose := false
defer func() {
rc.connPutChan <- &ConnPutTask{conn: conn, forceClose: forceClose}
rc.EnqueueConnTask(&ConnPutTask{conn: conn, forceClose: forceClose})
parts := strings.Split(addr, ":")
if len(parts) > 0 && addr != "" {
stat.EndStat(fmt.Sprintf("flashPutBlock:%v", parts[0]), err, bgTime, 1)
@ -957,7 +981,7 @@ func (rc *RemoteCacheClient) deleteRemoteBlock(key string, addr string) {
return
}
defer func() {
rc.connPutChan <- &ConnPutTask{conn: conn, forceClose: err != nil}
rc.EnqueueConnTask(&ConnPutTask{conn: conn, forceClose: err != nil})
}()
p := proto.NewPacketReqID()
p.Data = ([]byte)(key)
@ -1034,7 +1058,7 @@ func (rc *RemoteCacheClient) Read(ctx context.Context, fg *FlashGroup, reqId int
bgTime := stat.BeginStat()
defer func() {
forceClose := err != nil && !proto.IsFlashNodeLimitError(err)
rc.connPutChan <- &ConnPutTask{conn: conn, forceClose: forceClose}
rc.EnqueueConnTask(&ConnPutTask{conn: conn, forceClose: forceClose})
if err != nil && strings.Contains(err.Error(), "timeout") {
err = proto.ErrorReadTimeout
}
@ -1068,7 +1092,7 @@ func (rc *RemoteCacheClient) Read(ctx context.Context, fg *FlashGroup, reqId int
if err = reqPacket.WriteToConn(conn); err != nil {
log.LogWarnf("FlashGroup Read: failed to write to addr(%v) err(%v) remoteCacheMultiRead(%v)", addr, err, rc.RemoteCacheMultiRead)
rc.connPutChan <- &ConnPutTask{conn: conn, forceClose: err != nil}
rc.EnqueueConnTask(&ConnPutTask{conn: conn, forceClose: err != nil})
moved = fg.moveToUnknownRank(addr, err, rc.FlashNodeTimeoutCount)
if rc.RemoteCacheMultiRead {
log.LogInfof("Retrying due to write to addr(%v) failure err(%v)", addr, err)
@ -1082,7 +1106,7 @@ func (rc *RemoteCacheClient) Read(ctx context.Context, fg *FlashGroup, reqId int
break
}
log.LogWarnf("FlashGroup Read: getReadReply from addr(%v) err(%v) remoteCacheMultiRead(%v)", addr, err, rc.RemoteCacheMultiRead)
rc.connPutChan <- &ConnPutTask{conn: conn, forceClose: err != nil}
rc.EnqueueConnTask(&ConnPutTask{conn: conn, forceClose: err != nil})
moved = fg.moveToUnknownRank(addr, err, rc.FlashNodeTimeoutCount)
if rc.RemoteCacheMultiRead {
log.LogInfof("Retrying due to getReadReply from addr(%v) failure err(%v)", addr, err)
@ -1123,7 +1147,7 @@ func (rc *RemoteCacheClient) Prepare(ctx context.Context, fg *FlashGroup, req *p
return
}
defer func() {
rc.connPutChan <- &ConnPutTask{conn: conn, forceClose: err != nil}
rc.EnqueueConnTask(&ConnPutTask{conn: conn, forceClose: err != nil})
}()
if err = reqPacket.WriteToConn(conn); err != nil {
@ -1459,7 +1483,7 @@ func (reader *RemoteCacheReader) read(p []byte) (n int, err error) {
atomic.AddInt64(&reader.alreadyReadLen, expectLen)
if atomic.LoadInt64(&reader.alreadyReadLen) >= reader.needReadLen {
reader.loadAll = true
reader.rc.connPutChan <- &ConnPutTask{conn: reader.conn, forceClose: false}
reader.rc.EnqueueConnTask(&ConnPutTask{conn: reader.conn, forceClose: false})
reader.conn = nil
}
return int(expectLen), nil
@ -1469,7 +1493,7 @@ func (reader *RemoteCacheReader) Read(p []byte) (n int, err error) {
defer func() {
if err != nil && err != io.EOF && !reader.closed {
reader.closed = true
reader.rc.connPutChan <- &ConnPutTask{conn: reader.conn, forceClose: true}
reader.rc.EnqueueConnTask(&ConnPutTask{conn: reader.conn, forceClose: true})
if !reader.directRead {
bytespool.Free(reader.buffer)
}
@ -1535,7 +1559,7 @@ func (reader *RemoteCacheReader) Close() error {
if !reader.closed {
log.LogDebugf("RemoteCacheReader Close, reqId(%v)", reader.reqID)
reader.closed = true
reader.rc.connPutChan <- &ConnPutTask{conn: reader.conn, forceClose: reader.currOffset < reader.endOffset}
reader.rc.EnqueueConnTask(&ConnPutTask{conn: reader.conn, forceClose: reader.currOffset < reader.endOffset})
if !reader.directRead {
bytespool.Free(reader.buffer)
}
@ -1569,7 +1593,7 @@ func (rc *RemoteCacheClient) executeReadOperation(op *ReadOperation, isLast bool
if err = op.ReqPacket.WriteToNoDeadLineConn(conn); err != nil {
log.LogWarnf("%v FlashGroup Read: failed to write to addr(%v) err(%v)",
op.LogPrefix, op.Addr, err)
rc.connPutChan <- &ConnPutTask{conn: conn, forceClose: err != nil}
rc.EnqueueConnTask(&ConnPutTask{conn: conn, forceClose: err != nil})
if isLast && atomic.CompareAndSwapInt32(&op.HasResult, 0, 1) {
op.Conn = nil
op.BlockDataSize = 0
@ -1591,7 +1615,7 @@ func (rc *RemoteCacheClient) executeReadOperation(op *ReadOperation, isLast bool
log.LogDebugf("%v ReadObject getReadObjectReply from(%v) failed, reply(%v) error(%v)",
op.LogPrefix, conn.RemoteAddr(), reply, err)
}
rc.connPutChan <- &ConnPutTask{conn: conn, forceClose: !openConn}
rc.EnqueueConnTask(&ConnPutTask{conn: conn, forceClose: !openConn})
if openConn && atomic.CompareAndSwapInt32(&op.HasResult, 0, 1) {
op.Conn = nil
op.BlockDataSize = 0
@ -1606,7 +1630,7 @@ func (rc *RemoteCacheClient) executeReadOperation(op *ReadOperation, isLast bool
op.Err = nil
atomic.StoreInt32(&op.EndLoop, 1)
} else {
rc.connPutChan <- &ConnPutTask{conn: conn, forceClose: true}
rc.EnqueueConnTask(&ConnPutTask{conn: conn, forceClose: true})
}
}

View File

@ -220,20 +220,20 @@ func (rq *ReadTaskQueue) batchSmallObjectRead(taskWithPacket *ReadTaskWithPacket
}
conn.SetWriteDeadline(time.Now().Add(time.Duration(connTimeOut) * time.Millisecond))
if err = taskWithPacket.packet.WriteToNoDeadLineConn(conn); err != nil {
rq.rc.connPutChan <- &ConnPutTask{conn: conn, forceClose: true}
rq.rc.EnqueueConnTask(&ConnPutTask{conn: conn, forceClose: true})
log.LogWarnf("batchSmallObjectRead: failed to write, firstPacktTime(%v) addr(%v) err(%v)", rq.rc.firstPacketTimeout, rq.flashAddr, err)
return
}
replyPacket = NewFlashCacheReply()
if err = replyPacket.ReadFromConnExt(conn, connTimeOut); err != nil {
rq.rc.connPutChan <- &ConnPutTask{conn: conn, forceClose: true}
rq.rc.EnqueueConnTask(&ConnPutTask{conn: conn, forceClose: true})
log.LogWarnf("batchSmallObjectRead: failed to read, firstPacktTime(%v) addr(%v) err(%v)", rq.rc.firstPacketTimeout, rq.flashAddr, err)
return
}
getResult = true
atomic.StoreInt32(&taskWithPacket.finish, 1)
rq.rc.connPutChan <- &ConnPutTask{conn: conn, forceClose: false}
rq.rc.EnqueueConnTask(&ConnPutTask{conn: conn, forceClose: false})
if replyPacket.ResultCode != proto.OpOk {
err = fmt.Errorf(string(replyPacket.Data))
log.LogWarnf("batchSmallObjectRead: ResultCode NOK, replyPacket(%v), addr(%v), ResultCode(%v) err(%v)", replyPacket, rq.flashAddr, replyPacket.ResultCode, err.Error())