From a59c4c71edaea11a553a0ebb6db61ea234486161 Mon Sep 17 00:00:00 2001 From: chihe Date: Tue, 16 Sep 2025 09:42:42 +0800 Subject: [PATCH] fix(client): read from leader when excuting ahead read close:#1000340945 Signed-off-by: chihe --- sdk/data/stream/extent_handler.go | 3 +- sdk/data/stream/stream_aheadread.go | 163 ++++++++++++---------------- sdk/data/stream/stream_reader.go | 4 +- sdk/data/stream/stream_writer.go | 28 +++-- util/buf/buffer_pool.go | 4 + 5 files changed, 94 insertions(+), 108 deletions(-) diff --git a/sdk/data/stream/extent_handler.go b/sdk/data/stream/extent_handler.go index d80ce92b4..382893dc1 100644 --- a/sdk/data/stream/extent_handler.go +++ b/sdk/data/stream/extent_handler.go @@ -256,7 +256,8 @@ func (eh *ExtentHandler) write(data []byte, offset, size int, direct bool) (ek * for total < size { if eh.packet == nil { eh.packet = NewWritePacket(eh.inode, offset+total, eh.storeMode) - log.LogDebugf("ExtentHandler write packet nil and new packet: eh(%v)", eh) + log.LogDebugf("ExtentHandler write packet nil and new packet: eh(%v) offset(%v)", + eh, eh.packet.KernelOffset) if direct { eh.packet.Opcode = proto.OpSyncWrite } diff --git a/sdk/data/stream/stream_aheadread.go b/sdk/data/stream/stream_aheadread.go index 38352115b..d711d0e0d 100644 --- a/sdk/data/stream/stream_aheadread.go +++ b/sdk/data/stream/stream_aheadread.go @@ -18,7 +18,6 @@ import ( "fmt" "net" "runtime/debug" - "strings" "sync" "sync/atomic" "time" @@ -76,14 +75,15 @@ func putAheadReadBlock(block *AheadReadBlock) { } type AheadReadTask struct { - p *Packet - dnHosts []string - time time.Time - req *ExtentRequest - cacheSize int - cacheType string - logTime *time.Time - reqID string + p *Packet + dp *wrapper.DataPartition + time time.Time + req *ExtentRequest + cacheSize int + cacheType string + logTime *time.Time + reqID string + storageClass uint32 } type AheadReadWindow struct { @@ -230,8 +230,8 @@ func (arw *AheadReadWindow) backgroundAheadReadTask() { func (arw *AheadReadWindow) doTask(task *AheadReadTask) { var ( err error - hosts = task.dnHosts readBytes int + realHost = "invaild" ) key := fmt.Sprintf("%v-%v-%v-%v", task.p.inode, task.p.PartitionID, task.p.ExtentID, task.p.ExtentOffset) if arw.putTaskIfNotExist(key) { @@ -245,6 +245,7 @@ func (arw *AheadReadWindow) doTask(task *AheadReadTask) { } stat.EndStat("PrepareAheadData", err, task.logTime, 1) stat.EndStat(fmt.Sprintf("PrepareAheadData[%v]", task.cacheType), err, task.logTime, 1) + stat.EndStat(fmt.Sprintf("PrepareAheadData[%v][%v]", task.cacheType, realHost), err, task.logTime, 1) }() if _, ok := arw.cache.blockCache.Load(key); ok { return @@ -274,46 +275,48 @@ func (arw *AheadReadWindow) doTask(task *AheadReadTask) { cacheBlock.time = time.Now().Unix() atomic.StoreUint64(&cacheBlock.readBytes, 0) - for _, host := range hosts { - err = sendToNode(host, task.p, func(conn *net.TCPConn) (error, bool) { - // reset readBytes when reading from other datanodes - readBytes = 0 - atomic.StoreUint64(&cacheBlock.readBytes, 0) - for readBytes < task.cacheSize { - rp := NewReply(task.p.ReqID, task.p.PartitionID, task.p.ExtentID) - bufSize := util.Min(util.CacheReadBlockSize, task.cacheSize-readBytes) - rp.Data = cacheBlock.data[readBytes : readBytes+bufSize] - if e := rp.readFromConn(conn, proto.ReadDeadlineTime); e != nil { - return e, false - } - if rp.ResultCode != proto.OpOk { - return fmt.Errorf("resultCode(%x) not ok", rp.ResultCode), false - } - // update timeStamp to prevent from deleted by timeout - cacheBlock.time = time.Now().Unix() - readBytes += int(rp.Size) - atomic.AddUint64(&cacheBlock.readBytes, uint64(rp.Size)) - curSize := atomic.LoadUint64(&cacheBlock.readBytes) - log.LogDebugf("aheadRead doTask key(%v) curSize(%v) readBytes(%v) rp.Size(%v) task.cacheSize(%v) addr(%p)", - key, curSize, readBytes, rp.Size, task.cacheSize, cacheBlock) - if curSize > util.CacheReadBlockSize { - log.LogErrorf("aheadRead doTask out of range key(%v) curSize(%v) readBytes(%v) rp.Size(%v) task.cacheSize(%v) addr(%p)", - key, curSize, readBytes, rp.Size, task.cacheSize, cacheBlock) - } + enableFollowerRead := arw.streamer.client.dataWrapper.FollowerRead() && !arw.streamer.client.dataWrapper.InnerReq() + retryRead := true + if proto.IsCold(arw.streamer.client.volumeType) || proto.IsStorageClassBlobStore(task.storageClass) { + retryRead = false + } + sc := NewStreamConn(task.dp, enableFollowerRead, arw.streamer.client.streamRetryTimeout) + err = sc.Send(&retryRead, task.p, func(conn *net.TCPConn) (error, bool) { + // reset readBytes when reading from other datanodes + readBytes = 0 + atomic.StoreUint64(&cacheBlock.readBytes, 0) + for readBytes < task.cacheSize { + rp := NewReply(task.p.ReqID, task.p.PartitionID, task.p.ExtentID) + bufSize := util.Min(util.CacheReadBlockSize, task.cacheSize-readBytes) + rp.Data = cacheBlock.data[readBytes : readBytes+bufSize] + if e := rp.readFromConn(conn, proto.ReadDeadlineTime); e != nil { + return e, false + } + if rp.ResultCode != proto.OpOk { + return fmt.Errorf("resultCode(%x) not ok", rp.ResultCode), false + } + // update timeStamp to prevent from deleted by timeout + cacheBlock.time = time.Now().Unix() + readBytes += int(rp.Size) + atomic.AddUint64(&cacheBlock.readBytes, uint64(rp.Size)) + curSize := atomic.LoadUint64(&cacheBlock.readBytes) + log.LogDebugf("aheadRead doTask key(%v) curSize(%v) readBytes(%v) rp.Size(%v) task.cacheSize(%v) addr(%p)", + key, curSize, readBytes, rp.Size, task.cacheSize, cacheBlock) + if curSize > util.CacheReadBlockSize { + log.LogErrorf("aheadRead doTask out of range key(%v) curSize(%v) readBytes(%v) rp.Size(%v) task.cacheSize(%v) addr(%p)", + key, curSize, readBytes, rp.Size, task.cacheSize, cacheBlock) } - return nil, false - }) - if err != nil && strings.Contains(err.Error(), "timeout") { - log.LogWarnf("doTask try other host key(%v) addr(%p) read failed: err(%v)", key, cacheBlock, err) - continue - } - if err == nil { - atomic.StoreUint32(&cacheBlock.state, AheadReadBlockStateInit) - arw.canAheadRead = true - log.LogDebugf("aheadRead ready: key(%v) offset(%v) size(%v) err(%v) %v reqID(%v)", - task.req, task.p.ExtentOffset, task.cacheSize, err, time.Since(task.time), task.reqID) - return } + realHost = conn.RemoteAddr().String() + return nil, false + }) + + if err == nil { + atomic.StoreUint32(&cacheBlock.state, AheadReadBlockStateInit) + arw.canAheadRead = true + log.LogDebugf("aheadRead ready: key(%v) offset(%v) size(%v) err(%v) %v reqID(%v)", + task.req, task.p.ExtentOffset, task.cacheSize, err, time.Since(task.time), task.reqID) + return } atomic.StoreUint32(&cacheBlock.state, AheadReadBlockStateInit) log.LogErrorf("aheadRead doTask inode(%v) offset(%v) size(%v) err(%v) - read failed,"+ @@ -339,15 +342,16 @@ func (arw *AheadReadWindow) deleteTask(key string) { delete(arw.curTaskMap, key) } -func (arw *AheadReadWindow) addNextTask(offset int, dnHosts []string, req *ExtentRequest, stTime time.Time, reqID string) { +func (arw *AheadReadWindow) addNextTask(offset int, dp *wrapper.DataPartition, req *ExtentRequest, stTime time.Time, + reqID string, storageClass uint32) { id := offset/util.CacheReadBlockSize + arw.cache.winCnt remainSize := int(req.ExtentKey.Size) - id*util.CacheReadBlockSize if remainSize <= 0 { - arw.doMultiAheadRead(0, 0, req, dnHosts, stTime, reqID) + arw.doMultiAheadRead(0, 0, req, dp, stTime, reqID, storageClass) return } // log.LogDebugf("addNextTask ek(%v) remainSize(%v) %v-%v-%v-%v", req.ExtentKey, arw.streamer.inode, req.ExtentKey.PartitionId, req.ExtentKey.ExtentId, id*util.CacheReadBlockSize) - task := arw.getAheadReadTask(dnHosts, req, id, util.Min(remainSize, util.CacheReadBlockSize)) + task := arw.getAheadReadTask(dp, req, id, util.Min(remainSize, util.CacheReadBlockSize), storageClass) if task == nil { return } @@ -358,7 +362,8 @@ func (arw *AheadReadWindow) addNextTask(offset int, dnHosts []string, req *Exten arw.taskC <- task } -func (arw *AheadReadWindow) doMultiAheadRead(offset, remainSize int, req *ExtentRequest, dnHosts []string, startTime time.Time, reqID string) { +func (arw *AheadReadWindow) doMultiAheadRead(offset, remainSize int, req *ExtentRequest, dp *wrapper.DataPartition, + startTime time.Time, reqID string, storageClass uint32) { if len(arw.curTaskMap) > 0 && !arw.canAheadRead { return } @@ -372,10 +377,7 @@ func (arw *AheadReadWindow) doMultiAheadRead(offset, remainSize int, req *Extent for i := offset / util.CacheReadBlockSize; i < winCnt; i++ { if remainSize <= 0 { id = 0 - var ( - dp *wrapper.DataPartition - err error - ) + var err error curReq.ExtentKey = arw.streamer.getNextExtent(int(curReq.ExtentKey.FileOffset)) if curReq.ExtentKey == nil { return @@ -383,12 +385,11 @@ func (arw *AheadReadWindow) doMultiAheadRead(offset, remainSize int, req *Extent if dp, err = arw.streamer.client.dataWrapper.GetDataPartition(curReq.ExtentKey.PartitionId); err != nil { continue } - dnHosts = dp.Hosts remainSize = int(curReq.ExtentKey.Size) } size := util.Min(remainSize, util.CacheReadBlockSize) - task := arw.getAheadReadTask(dnHosts, curReq, id, size) + task := arw.getAheadReadTask(dp, curReq, id, size, storageClass) if task != nil { task.time = startTime task.cacheType = "pass" @@ -403,14 +404,15 @@ func (arw *AheadReadWindow) doMultiAheadRead(offset, remainSize int, req *Extent } } -func (arw *AheadReadWindow) getAheadReadTask(dnHosts []string, req *ExtentRequest, id int, size int) *AheadReadTask { +func (arw *AheadReadWindow) getAheadReadTask(dp *wrapper.DataPartition, req *ExtentRequest, id int, size int, storageClass uint32) *AheadReadTask { cacheOffset := int(id) * util.CacheReadBlockSize p := NewReadPacket(req.ExtentKey, cacheOffset, size, arw.streamer.inode, req.FileOffset, true) task := &AheadReadTask{ - p: p, - dnHosts: getRandomHostById(dnHosts, id), - req: req, - cacheSize: size, + p: p, + dp: dp, + req: req, + cacheSize: size, + storageClass: storageClass, } return task } @@ -449,7 +451,7 @@ func (s *Streamer) getNextExtent(offset int) (ek *proto.ExtentKey) { return } -func (s *Streamer) aheadRead(req *ExtentRequest) (readSize int, err error) { +func (s *Streamer) aheadRead(req *ExtentRequest, storageClass uint32) (readSize int, err error) { var ( offset int cacheOffset int @@ -511,7 +513,7 @@ func (s *Streamer) aheadRead(req *ExtentRequest) (readSize int, err error) { log.LogDebugf("aheadRead move ahead win inode(%v) FileOffset(%v) offset(%v) "+ "size(%v) cacheBlockOffset(%v) cacheBlockSize(%v) reqID(%v)", s.inode, req.FileOffset, offset, needSize, cacheBlock.offset, curSize, reqID) - go s.aheadReadWindow.addNextTask(offset, dp.Hosts, req, startTime, reqID) + go s.aheadReadWindow.addNextTask(offset, dp, req, startTime, reqID, storageClass) copy(req.Data, cacheBlock.data[offset-int(cacheBlock.offset):int(curSize)]) cacheBlock.lock.RUnlock() readSize += int(curSize) + int(cacheBlock.offset) - offset @@ -528,35 +530,6 @@ func (s *Streamer) aheadRead(req *ExtentRequest) (readSize int, err error) { break } step = "pass" - go s.aheadReadWindow.doMultiAheadRead(offset, remainSize, req, dp.Hosts, startTime, uuid.New().String()) + go s.aheadReadWindow.doMultiAheadRead(offset, remainSize, req, dp, startTime, uuid.New().String(), storageClass) return } - -func sendToNode(host string, p *Packet, getReply GetReplyFunc) (err error) { - var conn *net.TCPConn - if conn, err = StreamConnPool.GetConnect(host); err != nil { - return - } - defer StreamConnPool.PutConnect(conn, err != nil) - if err = p.WriteToConn(conn); err != nil { - return - } - err, _ = getReply(conn) - return -} - -func getRandomHostById(hosts []string, id int) []string { - if len(hosts) == 0 { - return nil - } - hIdx := id % len(hosts) - rhosts := make([]string, len(hosts)) - for i := 0; i < len(hosts); i++ { - if hIdx >= len(hosts) { - hIdx = 0 - } - rhosts[i] = hosts[hIdx] - hIdx++ - } - return rhosts -} diff --git a/sdk/data/stream/stream_reader.go b/sdk/data/stream/stream_reader.go index 4070ace83..2ff2f5750 100644 --- a/sdk/data/stream/stream_reader.go +++ b/sdk/data/stream/stream_reader.go @@ -258,7 +258,7 @@ func (s *Streamer) read(data []byte, offset int, size int, storageClass uint32) } else { log.LogDebugf("Stream read: ino(%v) req(%v) s.needBCache(%v) s.client.bcacheEnable(%v) aheadReadEnable(%v)", s.inode, req, s.needBCache, s.client.bcacheEnable, s.aheadReadEnable) if s.aheadReadEnable && req.ExtentKey.Size > util.CacheReadBlockSize { - readBytes, err = s.aheadRead(req) + readBytes, err = s.aheadRead(req, storageClass) if err == nil && readBytes == req.Size { total += readBytes continue @@ -579,7 +579,7 @@ func (s *Streamer) completeAsyncFlush(req *AsyncFlushRequest) { // the flush order cannot be guaranteed err := handler.flush() if err != nil { - log.LogWarnf("completeAsyncFlush: completed successfully for handler(%v)", handler) + log.LogWarnf("completeAsyncFlush: completed failed for handler(%v)", handler) } else { log.LogDebugf("completeAsyncFlush: completed successfully for handler(%v) err(%v)", handler, err) if req.clearFunc != nil { diff --git a/sdk/data/stream/stream_writer.go b/sdk/data/stream/stream_writer.go index 17ee8478d..911c73991 100644 --- a/sdk/data/stream/stream_writer.go +++ b/sdk/data/stream/stream_writer.go @@ -939,7 +939,8 @@ func (s *Streamer) doWriteAppendEx(data []byte, offset, size int, direct bool, r break } - log.LogDebugf("doWriteAppendEx: start write protection, handler(%v)", s.handler) + log.LogDebugf("doWriteAppendEx: start write protection, handler(%v) offset(%v) size(%v)", + s.handler, offset, size) ek, err = s.handler.write(data, offset, size, direct) if err == nil && ek != nil { ek.SetSeq(s.verSeq) @@ -1005,7 +1006,7 @@ func (s *Streamer) doWriteAppendEx(data []byte, offset, size int, direct bool, r } func (s *Streamer) flush(wait bool, id string) (err error) { - pending := make(map[uint64]chan error) + pending := make(map[*ExtentHandler]chan error) elements := s.dirtylist.Elements() for _, element := range elements { if element == nil { @@ -1029,8 +1030,8 @@ func (s *Streamer) flush(wait bool, id string) (err error) { s.inode, eh, atomic.LoadInt32(&eh.inflight), id, wait) ch := s.requestAsyncFlush(eh, clearFunc) if wait { - if _, ok := pending[eh.id]; !ok { - pending[eh.id] = ch + if _, ok := pending[eh]; !ok { + pending[eh] = ch } eh.stream.dirtylist.Remove(element) } @@ -1053,15 +1054,19 @@ func (s *Streamer) flush(wait bool, id string) (err error) { } log.LogDebugf("Streamer(%v) wait(%v) pending(%v) id(%v)", s.inode, wait, pending, id) if wait { + var firstErr error for len(pending) > 0 { progressed := false - for hid, ch := range pending { + for eh, ch := range pending { select { - case err = <-ch: - s.removePendingAsyncFlush(hid) - delete(pending, hid) - if err != nil { - return err + case e := <-ch: + s.removePendingAsyncFlush(eh.id) + delete(pending, eh) + if e != nil && firstErr == nil { + firstErr = e + } + if e != nil { + eh.cleanup() } progressed = true default: @@ -1072,6 +1077,9 @@ func (s *Streamer) flush(wait bool, id string) (err error) { time.Sleep(time.Millisecond * 10) } } + if firstErr != nil { + return firstErr + } } return } diff --git a/util/buf/buffer_pool.go b/util/buf/buffer_pool.go index 175c7580f..3beebce35 100644 --- a/util/buf/buffer_pool.go +++ b/util/buf/buffer_pool.go @@ -5,8 +5,10 @@ import ( "fmt" "sync" "sync/atomic" + "time" "github.com/cubefs/cubefs/util" + "github.com/cubefs/cubefs/util/log" "golang.org/x/time/rate" ) @@ -236,7 +238,9 @@ func (bufferP *BufferPool) getHeadProtoVer(id uint64) (data []byte) { func (bufferP *BufferPool) getNormal(id uint64) (data []byte) { if NormalBuffersTotalLimit != InvalidLimit && atomic.LoadInt64(&normalBuffersCount) >= NormalBuffersTotalLimit { ctx := context.Background() + start := time.Now() normalBuffersRateLimit.Wait(ctx) + log.LogDebugf("getNormal wait: id(%v) (%v)", id, time.Since(start).String()) } select { case data = <-bufferP.normalPools[id%slotCnt]: