mirror of
https://github.com/cubefs/cubefs.git
synced 2026-08-02 02:00:56 +00:00
fix(flashnode): clean streamer resource at consumer finish
with: #1000150848
Signed-off-by: clinx <chenlin1@oppo.com>
(cherry picked from commit c3737e9b38)
This commit is contained in:
parent
0270edc70e
commit
2ce4af6a64
@ -48,7 +48,7 @@ const (
|
||||
DefaultExpireTime = 60 * 60
|
||||
InitFileName = "flash.init"
|
||||
DefaultCacheDirName = "cache"
|
||||
DefaultCacheMaxUsedRatio = 0.95
|
||||
DefaultCacheMaxUsedRatio = 0.99
|
||||
DefaultEnableTmpfs = true
|
||||
|
||||
LRUCacheBlockCacheType = 0
|
||||
|
||||
@ -70,7 +70,7 @@ const (
|
||||
_defaultHandlerFileRoutineNumPerTask = 20
|
||||
_maxHandlerFileRoutineNumPerTask = 500
|
||||
_defaultManualScanLimitPerSecond = 10000
|
||||
_defaultPrepareLimitPerSecond = 1000
|
||||
_defaultPrepareLimitPerSecond = 10000
|
||||
_defaultManualScanLimitBurst = 1000
|
||||
_slotStatValidPeriod = 10 * time.Minute // min
|
||||
_defaultPrepareRoutineNum = 20
|
||||
|
||||
@ -317,21 +317,24 @@ func (s *ManualScanner) warmUp(i *proto.ScanItem) error {
|
||||
log.LogInfof("skip warmUp, len of extents=0, inode(%v)", i.Inode)
|
||||
return nil
|
||||
}
|
||||
defer func() {
|
||||
if err != nil {
|
||||
s.ec.CloseStream(i.Inode)
|
||||
s.ec.EvictStream(i.Inode)
|
||||
}
|
||||
}()
|
||||
if err = s.ec.OpenStream(i.Inode, false, false, ""); err != nil {
|
||||
log.LogWarnf("warmUp: ec OpenStream fail, inode(%v) err: %v", i.Inode, err)
|
||||
return err
|
||||
}
|
||||
defer func() {
|
||||
s.ec.CloseStream(i.Inode)
|
||||
s.ec.EvictStream(i.Inode)
|
||||
}()
|
||||
if err = s.ec.ForceRefreshExtentsCache(i.Inode); err != nil {
|
||||
log.LogWarnf("warmUp: ec ForceRefreshExtentsCache fail, inode(%v) err: %v", i.Inode, err)
|
||||
return err
|
||||
}
|
||||
for _, extent := range extents {
|
||||
eLen := len(extents)
|
||||
for index, extent := range extents {
|
||||
s.prepareLimiter.Wait(context.Background())
|
||||
prepareReq := stream.NewPrepareRemoteCacheRequest(i.Inode, extent, true, i.WriteGen)
|
||||
prepareReq := stream.NewPrepareRemoteCacheRequest(i.Inode, extent, true, i.WriteGen, eLen-1 == index)
|
||||
s.RemoteCache.PrepareCh <- prepareReq
|
||||
atomic.AddInt64(&s.currentStat.TotalExtentKeyNum, 1)
|
||||
atomic.AddInt64(&s.currentStat.TotalCacheSize, int64(extent.Size))
|
||||
|
||||
@ -969,6 +969,10 @@ func (c *ExtentClient) servePrepareRequest(prepareReq *PrepareRemoteCacheRequest
|
||||
}
|
||||
if prepareReq.warmUp {
|
||||
s.prepareRemoteCache(prepareReq.ctx, prepareReq.ek, prepareReq.gen)
|
||||
if prepareReq.triggerClean {
|
||||
s.client.CloseStream(prepareReq.inode)
|
||||
s.client.EvictStream(prepareReq.inode)
|
||||
}
|
||||
} else {
|
||||
inodeInfo, err := s.client.getInodeInfo(prepareReq.inode)
|
||||
if err != nil {
|
||||
|
||||
@ -29,20 +29,22 @@ import (
|
||||
const SIZE_GB = 1024 * 1024 * 1024
|
||||
|
||||
type PrepareRemoteCacheRequest struct {
|
||||
ctx context.Context
|
||||
inode uint64
|
||||
ek *proto.ExtentKey
|
||||
warmUp bool
|
||||
gen uint64
|
||||
ctx context.Context
|
||||
inode uint64
|
||||
ek *proto.ExtentKey
|
||||
warmUp bool
|
||||
gen uint64
|
||||
triggerClean bool
|
||||
}
|
||||
|
||||
func NewPrepareRemoteCacheRequest(inode uint64, ek proto.ExtentKey, warmUp bool, gen uint64) *PrepareRemoteCacheRequest {
|
||||
func NewPrepareRemoteCacheRequest(inode uint64, ek proto.ExtentKey, warmUp bool, gen uint64, triggerClean bool) *PrepareRemoteCacheRequest {
|
||||
return &PrepareRemoteCacheRequest{
|
||||
ctx: context.Background(),
|
||||
inode: inode,
|
||||
ek: &ek,
|
||||
warmUp: warmUp,
|
||||
gen: gen,
|
||||
ctx: context.Background(),
|
||||
inode: inode,
|
||||
ek: &ek,
|
||||
warmUp: warmUp,
|
||||
gen: gen,
|
||||
triggerClean: triggerClean,
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
Loading…
Reference in New Issue
Block a user