From 9c22ca99ca85bbbca987e4dd5279265f47f2003a Mon Sep 17 00:00:00 2001 From: chihe Date: Fri, 31 Oct 2025 15:28:10 +0800 Subject: [PATCH] fix(client): use extentOffset to create block key for aheadRead close:#1000412030 Signed-off-by: chihe --- sdk/data/stream/stream_aheadread.go | 24 +++++++++++++++--------- 1 file changed, 15 insertions(+), 9 deletions(-) diff --git a/sdk/data/stream/stream_aheadread.go b/sdk/data/stream/stream_aheadread.go index 508322638..4dd4a75c1 100644 --- a/sdk/data/stream/stream_aheadread.go +++ b/sdk/data/stream/stream_aheadread.go @@ -186,7 +186,7 @@ func (arc *AheadReadCache) doCheckBlockTimeOut() { if curTime-bv.time > int64(arc.blockTimeOut) { log.LogDebugf("doCheckBlockTimeOut delete readed block: key(%v) addr(%p)", key, bv) if atomic.LoadUint32(&bv.readed) == AheadReadedInit { - log.LogWarnf("doCheckBlockTimeOut delete unreaded block: key(%v) addr(%p)", key, bv) + log.LogInfof("doCheckBlockTimeOut delete unreaded block: key(%v) addr(%p)", key, bv) } arc.blockCache.Delete(key.(string)) // block is ready to recycle @@ -244,9 +244,9 @@ func (arw *AheadReadWindow) doTask(task *AheadReadTask) { key string ) if readDataFromTinyExtent(task.p.ExtentID) { - key = fmt.Sprintf("%v-%v-%v-%v", task.p.inode, task.p.PartitionID, task.p.ExtentID, 0) + key = createAheadBlockKey(task.p.inode, task.p.PartitionID, task.p.ExtentID, uint64(task.p.ExtentOffset), 0) } else { - key = fmt.Sprintf("%v-%v-%v-%v", task.p.inode, task.p.PartitionID, task.p.ExtentID, task.p.ExtentOffset) + key = createAheadBlockKey(task.p.inode, task.p.PartitionID, task.p.ExtentID, 0, int(task.p.ExtentOffset)) } if _, ok := arw.cache.blockCache.Load(key); ok { return @@ -368,7 +368,7 @@ func (arw *AheadReadWindow) addNextTask(offset int, stTime time.Time, reqID stri rightOffset := atomic.LoadUint64(&arw.rightOffset) ek := arw.streamer.getCurrentExtent(int(rightOffset)) if ek == nil { - log.LogWarnf("addNextTask current ExtentKey for offset(%v) is nil reqID(%v)", rightOffset, reqID) + log.LogInfof("addNextTask current ExtentKey for offset(%v) is nil reqID(%v)", rightOffset, reqID) return } // move forward @@ -413,7 +413,7 @@ func (arw *AheadReadWindow) doMultiAheadRead(offset int, req *ExtentRequest, dp offset = 0 } cacheOffset := offset / int(arw.streamer.aheadReadBlockSize) * int(arw.streamer.aheadReadBlockSize) - key := fmt.Sprintf("%v-%v-%v-%v", arw.streamer.inode, req.ExtentKey.PartitionId, req.ExtentKey.ExtentId, cacheOffset) + key := createAheadBlockKey(arw.streamer.inode, req.ExtentKey.PartitionId, req.ExtentKey.ExtentId, req.ExtentKey.ExtentOffset, cacheOffset) log.LogDebugf("doMultiAheadRead send: key(%v) reqID(%v) index(%v) req(%v)", key, reqID, offset/int(arw.streamer.aheadReadBlockSize), req) id := offset / int(arw.streamer.aheadReadBlockSize) @@ -445,7 +445,7 @@ func (arw *AheadReadWindow) doMultiAheadRead(offset int, req *ExtentRequest, dp task.cacheType = "pass" task.logTime = stat.BeginStat() task.reqID = reqID - key := fmt.Sprintf("%v-%v-%v-%v", arw.streamer.inode, task.p.PartitionID, task.p.ExtentID, task.p.ExtentOffset) + key := createAheadBlockKey(arw.streamer.inode, task.p.PartitionID, task.p.ExtentID, uint64(task.p.ExtentOffset), cacheOffset) log.LogDebugf("doMultiAheadRead send: key(%v) curReq(%v) size(%v) reqID(%v)", key, curReq, size, reqID) arw.taskC <- task @@ -459,7 +459,7 @@ func (arw *AheadReadWindow) getAheadReadTask(dp *wrapper.DataPartition, req *Ext ) *AheadReadTask { cacheOffset := id * int(arw.streamer.aheadReadBlockSize) if dp == nil { - key := fmt.Sprintf("%v-%v-%v-%v", arw.streamer.inode, req.ExtentKey.PartitionId, req.ExtentKey.ExtentId, cacheOffset) + key := createAheadBlockKey(arw.streamer.inode, req.ExtentKey.PartitionId, req.ExtentKey.ExtentId, req.ExtentKey.ExtentOffset, cacheOffset) log.LogWarnf("getAheadReadTask dp for key(%v) req(%v) is nil", key, req) return nil @@ -483,7 +483,7 @@ func (arw *AheadReadWindow) getAheadReadTask(dp *wrapper.DataPartition, req *Ext func (arw *AheadReadWindow) evictCacheBlock(req *ExtentRequest) { offset := req.FileOffset - int(req.ExtentKey.FileOffset) + int(req.ExtentKey.ExtentOffset) cacheOffset := offset / int(arw.streamer.aheadReadBlockSize) * int(arw.streamer.aheadReadBlockSize) - key := fmt.Sprintf("%v-%v-%v-%v", arw.streamer.inode, req.ExtentKey.PartitionId, req.ExtentKey.ExtentId, cacheOffset) + key := createAheadBlockKey(arw.streamer.inode, req.ExtentKey.PartitionId, req.ExtentKey.ExtentId, req.ExtentKey.ExtentOffset, cacheOffset) value, ok := arw.cache.blockCache.Load(key) if ok { bv := value.(*AheadReadBlock) @@ -541,7 +541,8 @@ func (s *Streamer) aheadRead(req *ExtentRequest, storageClass uint32) (readSize cacheOffset = offset / int(s.aheadReadBlockSize) * int(s.aheadReadBlockSize) needSize := req.Size - readSize for needSize > 0 { - key = fmt.Sprintf("%v-%v-%v-%v", s.inode, req.ExtentKey.PartitionId, req.ExtentKey.ExtentId, cacheOffset) + key = createAheadBlockKey(s.inode, req.ExtentKey.PartitionId, req.ExtentKey.ExtentId, + req.ExtentKey.ExtentOffset, cacheOffset) val, ok := s.aheadReadWindow.cache.blockCache.Load(key) if !ok { log.LogDebugf("aheadRead cache block key(%v) not found", key) @@ -691,3 +692,8 @@ func (s *Streamer) getCurrentExtent(offset int) (ek *proto.ExtentKey) { func readDataFromTinyExtent(extentID uint64) bool { return extentID <= 64 } + +func createAheadBlockKey(inode, partitionId, extentId, extentOffset uint64, cacheOffset int) string { + return fmt.Sprintf("%v-%v-%v-%v-%v", inode, partitionId, extentId, + extentOffset, cacheOffset) +}