mirror of
https://github.com/cubefs/cubefs.git
synced 2026-08-02 02:00:56 +00:00
fix(client): use extentOffset to create block key for aheadRead
close:#1000412030 Signed-off-by: chihe <chihe@oppo.com>
This commit is contained in:
parent
c6f85392b3
commit
9c22ca99ca
@ -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)
|
||||
}
|
||||
|
||||
Loading…
Reference in New Issue
Block a user