fix(flashnode): close cached files concurrently#23092015

Signed-off-by: clinx <chenlin1@oppo.com>
This commit is contained in:
clinx 2025-03-18 14:22:13 +08:00 committed by zhumingze1108
parent a1a2add9d9
commit a6939d0981
2 changed files with 31 additions and 10 deletions

View File

@ -405,15 +405,21 @@ func (c *CacheEngine) Start() (err error) {
func (c *CacheEngine) Stop() (err error) {
c.closeOnce.Do(func() { close(c.closeCh) })
var wg sync.WaitGroup
c.lruCacheMap.Range(func(key, value interface{}) bool {
wg.Add(1)
dataPath := key.(string)
cacheItem := value.(*lruCacheItem)
if err = cacheItem.lruCache.Close(); err != nil {
return true
}
log.LogInfof("CacheEngine stopped, data dir: %s", dataPath)
go func(d string, ci *lruCacheItem) {
defer wg.Done()
if err = cacheItem.lruCache.Close(); err != nil {
return
}
log.LogInfof("CacheEngine stopped, data dir: %s", dataPath)
}(dataPath, cacheItem)
return true
})
wg.Wait()
if err != nil {
return err
}

View File

@ -413,15 +413,30 @@ func (c *fCache) Close() error {
})
c.lock.Lock()
defer c.lock.Unlock()
var errCount int
// todo: split close action with big lock
chanItems := make(chan interface{}, 16)
var (
errCount int32
wg sync.WaitGroup
)
for i := 0; i < 8; i++ {
wg.Add(1)
go func() {
defer wg.Done()
for e := range chanItems {
err := c.onClose(e)
if err != nil {
atomic.AddInt32(&errCount, 1)
}
}
}()
}
for _, item := range c.items {
kv := item.Value.(*entry)
err := c.onClose(kv.value)
if err != nil {
errCount++
}
chanItems <- kv.value
}
close(chanItems)
wg.Wait()
if errCount > 0 {
return fmt.Errorf("error count(%v) on close", errCount)
}