fix(flashnode): rate limiting based on actual written data flow

close:#1000150848

Signed-off-by: clinx <chenlin1@oppo.com>

(cherry picked from commit 4320c3fc7c)
This commit is contained in:
clinx 2025-06-23 14:07:08 +08:00 committed by chihe
parent 23763d4219
commit d760b25ce1
5 changed files with 122 additions and 76 deletions

View File

@ -414,8 +414,13 @@ func testWriteSingleFileV2(t *testing.T) {
file, _ := cacheBlock.GetOrOpenFileHandler()
_, _ = file.ReadAt(crcBuf[:4], 8192+HeaderSize)
require.Equal(t, crcSum1, binary.BigEndian.Uint32(crcBuf[:4]))
fdata := make([]byte, 4096)
_, _ = file.ReadAt(fdata[:4096], HeaderSize)
require.Equal(t, crcSum1, crc32.ChecksumIEEE(fdata[:4096]))
_, _ = file.ReadAt(crcBuf[:4], 8192+HeaderSize+4)
require.Equal(t, crcSum2, binary.BigEndian.Uint32(crcBuf[:4]))
_, _ = file.ReadAt(fdata[:4096], HeaderSize+4096)
require.Equal(t, crcSum2, crc32.ChecksumIEEE(fdata[:4096]))
}
func testWriteSingleFileErrorV2(t *testing.T) {

View File

@ -320,75 +320,82 @@ func (f *FlashNode) opCachePutBlock(conn net.Conn, p *proto.Packet) (err error)
log.LogDebug("action[opCachePutBlock] create block key:"+uniKey+" logMsg:%s",
p.LogMessage(p.GetOpMsg(), conn.RemoteAddr().String(), p.StartT, err))
}
missTaskDone := make(chan struct{})
allocSize := cachengine.CalcAllocSizeV2(int(req.BlockLen))
err = f.limitWrite.TryRunAsync(context.Background(), allocSize, false, func() {
defer func() {
close(missTaskDone)
}()
if cb, err2, created := f.cacheEngine.CreateBlockV2(pDir, uniKey, req.TTL, uint32(allocSize), conn.RemoteAddr().String()); err2 != nil || created {
err = fmt.Errorf("already create block(%v)", cachengine.GenCacheBlockKeyV2(pDir, uniKey))
if cb, err2, created := f.cacheEngine.CreateBlockV2(pDir, uniKey, req.TTL, uint32(allocSize), conn.RemoteAddr().String()); err2 != nil || created {
err = fmt.Errorf("already create block(%v)", cachengine.GenCacheBlockKeyV2(pDir, uniKey))
return
} else {
p.PacketOkReply()
if err1 = p.WriteToConn(conn); err1 != nil {
log.LogErrorf("action[opCachePutBlock] write to conn %v", err1)
return
} else {
p.PacketOkReply()
if e := p.WriteToConn(conn); e != nil {
log.LogErrorf("action[opCachePutBlock] write to conn %v", e)
}
bgTime1 := stat.BeginStat()
defer func() {
stat.EndStat("CachePutBlock:Write", err1, bgTime1, 1)
}()
var (
totalWritten int64
n int
readSize int
)
buf := bytespool.Alloc(proto.PageSize)
defer bytespool.Free(buf)
crcBuf := bytespool.Alloc(cachengine.CRCLen)
defer bytespool.Free(crcBuf)
for {
readSize = proto.PageSize
}
bgTime1 := stat.BeginStat()
defer func() {
stat.EndStat("CachePutBlock:Write", err1, bgTime1, 1)
}()
var (
totalWritten int64
n int
readSize int
)
writeLen := proto.PageSize + 4
buf := bytespool.Alloc(writeLen)
defer bytespool.Free(buf)
for {
readSize = proto.PageSize
missTaskDone := make(chan struct{})
err = f.limitWrite.TryRunAsync(context.Background(), writeLen, false, func() {
defer func() {
close(missTaskDone)
}()
if totalWritten+int64(readSize) > req.BlockLen {
readSize = int(req.BlockLen - totalWritten)
}
if n, err1 = io.ReadFull(conn, buf[:proto.PageSize]); err1 != nil {
if n, err1 = io.ReadFull(conn, buf[:writeLen]); err1 != nil {
log.LogWarnf("action[opCachePutBlock] read data and crc from conn %v", err1)
return
}
if n != proto.PageSize {
if n != writeLen {
err1 = syscall.EBADMSG
}
if n, err1 = io.ReadFull(conn, crcBuf[:cachengine.CRCLen]); err1 != nil {
return
}
if n != cachengine.CRCLen {
err1 = syscall.EBADMSG
}
err1 = cb.WriteAtV2(&proto.FlashWriteParam{
Offset: totalWritten,
Size: req.BlockLen,
Data: buf[:proto.PageSize],
Crc: crcBuf[:cachengine.CRCLen],
Crc: buf[proto.PageSize:writeLen],
DataSize: int64(readSize),
})
if err1 != nil {
return
}
cachengine.UpdateWriteBytesMetric(proto.PageSize, cb.GetRootPath())
cachengine.UpdateWriteCountMetric(cb.GetRootPath())
err1 = cb.MaybeWriteCompleted(req.BlockLen)
if err1 != nil {
return
}
totalWritten += int64(readSize)
if totalWritten == req.BlockLen {
if log.EnableDebug() {
log.LogDebugf("action[opCachePutBlock] total write %v", totalWritten)
}
})
if err == nil {
<-missTaskDone
}
if err1 != nil || err != nil {
break
} else {
if err1 = p.WriteToConn(conn); err1 != nil {
log.LogErrorf("action[opCachePutBlock] reply to conn %v for write data", err1)
break
}
}
totalWritten += int64(readSize)
if totalWritten == req.BlockLen {
if log.EnableDebug() {
log.LogDebugf("action[opCachePutBlock] total write %v", totalWritten)
}
break
}
}
})
if err == nil {
<-missTaskDone
}
if err1 != nil {
f.cacheEngine.DeleteCacheBlock(blockKey)

View File

@ -197,13 +197,10 @@ func testTCPCachePutBlock(t *testing.T) {
require.NoError(t, p.WriteToConn(conn))
require.NoError(t, r.ReadFromConn(conn, 3))
require.Equal(t, proto.OpOk, r.ResultCode)
buf := make([]byte, proto.PageSize)
buf := make([]byte, proto.PageSize+4)
buf[0] = '{'
_, err := conn.Write(buf[:proto.PageSize])
require.NoError(t, err)
crcBuf := make([]byte, 4)
binary.BigEndian.PutUint32(crcBuf, crc32.ChecksumIEEE(buf[:proto.PageSize]))
_, err = conn.Write(crcBuf[:4])
binary.BigEndian.PutUint32(buf[proto.PageSize:], crc32.ChecksumIEEE(buf[:proto.PageSize]))
_, err := conn.Write(buf[:proto.PageSize+4])
require.NoError(t, err)
}

View File

@ -52,14 +52,15 @@ type AddressPingStats struct {
}
type RemoteCacheClient struct {
flashGroups *btree.BTree
mc *master.MasterClient
conns *util.ConnectPool
hostLatency sync.Map
TTL int64
ReadTimeout int64 // ms
stopC chan struct{}
wg sync.WaitGroup
flashGroups *btree.BTree
mc *master.MasterClient
conns *util.ConnectPool
hostLatency sync.Map
TTL int64
ReadTimeout int64 // ms
WriteTimeout int64
stopC chan struct{}
wg sync.WaitGroup
blockSize uint64
clusterEnabled bool
@ -101,6 +102,7 @@ func NewRemoteCacheClient(masters []string, blockSize uint64) (rc *RemoteCacheCl
rc.stopC = make(chan struct{})
rc.flashGroups = btree.New(32)
rc.ReadTimeout = proto.DefaultRemoteCacheClientReadTimeout
rc.WriteTimeout = proto.DefaultRemoteCacheExtentReadTimeout
rc.mc = master.NewMasterClient(masters, false)
rc.conns = util.NewConnectPoolWithTimeoutAndCap(5, 500, ConnIdelTimeout, 1)
@ -430,7 +432,7 @@ func (rc *RemoteCacheClient) Put(ctx context.Context, reqId, key string, r io.Re
addr := fg.getFlashHost()
if addr == "" {
err = fmt.Errorf("getFlashHost failed: can not find host")
log.LogWarnf("FlashGroup reqId(%v) put failed: err(%v)", reqId, err)
log.LogWarnf("FlashGroup fg(%v) reqId(%v) put failed: err(%v)", fg, reqId, err)
return
}
var conn *net.TCPConn
@ -441,7 +443,8 @@ func (rc *RemoteCacheClient) Put(ctx context.Context, reqId, key string, r io.Re
}
bgTime := stat.BeginStat()
defer func() {
rc.conns.PutConnect(conn, err != nil)
forceClose := err != nil && !proto.IsFlashNodeLimitError(err)
rc.conns.PutConnect(conn, forceClose)
parts := strings.Split(addr, ":")
if len(parts) > 0 && addr != "" {
stat.EndStat(fmt.Sprintf("flashPutBlock:%v", parts[0]), err, bgTime, 1)
@ -461,23 +464,30 @@ func (rc *RemoteCacheClient) Put(ctx context.Context, reqId, key string, r io.Re
return
}
replyPacket := NewFlashCacheReply()
if err = replyPacket.ReadFromConnExt(conn, int(rc.ReadTimeout)); err != nil {
if err = replyPacket.ReadFromConnExt(conn, int(rc.WriteTimeout)); err != nil {
log.LogWarnf("FlashGroup put: reqId(%v) failed to ReadFromConn, replyPacket(%v), fg host(%v) , err(%v)", reqId, replyPacket, addr, err)
return
}
if replyPacket.ResultCode != proto.OpOk {
err = fmt.Errorf("%v", string(replyPacket.Data))
if !proto.IsFlashNodeLimitError(err) {
log.LogWarnf("getPutBlockReply: ResultCode NOK, req(%v) reply(%v) ResultCode(%v)", reqPacket, replyPacket, replyPacket.ResultCode)
}
return
}
defer func() {
if err != nil {
rc.deleteRemoteBlock(key, conn)
log.LogWarnf("FlashGroup put: reqId(%v) remove key %v by err %v", reqId, key, err)
rc.deleteRemoteBlock(key, addr)
}
}()
var (
totalWritten int64
n int
)
buf := bytespool.Alloc(proto.PageSize)
writeLen := proto.PageSize + 4
buf := bytespool.Alloc(writeLen)
defer bytespool.Free(buf)
crcBuf := bytespool.Alloc(4)
defer bytespool.Free(crcBuf)
for {
if ctx.Err() != nil {
return ctx.Err()
@ -488,18 +498,29 @@ func (rc *RemoteCacheClient) Put(ctx context.Context, reqId, key string, r io.Re
}
n, err = r.Read(buf[:readSize])
if n != readSize {
log.LogWarnf("FlashGroup put: expected to read %d bytes, but only read %d", readSize, n)
return fmt.Errorf("expected to read %d bytes, but only read %d", readSize, n)
}
if n > 0 {
if log.EnableDebug() {
log.LogDebugf("FlashGroup put: write %d bytes total(%d) to fg", n, totalWritten)
}
binary.BigEndian.PutUint32(buf[proto.PageSize:], crc32.ChecksumIEEE(buf[:proto.PageSize]))
conn.SetWriteDeadline(time.Now().Add(proto.WriteDeadlineTime * time.Second))
if _, err = conn.Write(buf[:proto.PageSize]); err != nil {
log.LogErrorf("wirte data to flashnode get err %v", err)
if _, err = conn.Write(buf[:writeLen]); err != nil {
log.LogErrorf("wirte data and crc to flashnode get err %v", err)
return err
}
binary.BigEndian.PutUint32(crcBuf, crc32.ChecksumIEEE(buf[:proto.PageSize]))
if _, err = conn.Write(crcBuf[:4]); err != nil {
log.LogErrorf("wirte crc to flashnode get err %v", err)
return err
if err = replyPacket.ReadFromConnExt(conn, int(rc.WriteTimeout)); err != nil {
log.LogWarnf("FlashGroup put data: reqId(%v) failed to ReadFromConn, replyPacket(%v), fg host(%v) , err(%v)", reqId, replyPacket, addr, err)
return
}
if replyPacket.ResultCode != proto.OpOk {
err = fmt.Errorf("%v", string(replyPacket.Data))
if !proto.IsFlashNodeLimitError(err) {
log.LogWarnf("getPutBlockReply: put data ResultCode NOK, req(%v) reply(%v) ResultCode(%v)", reqPacket, replyPacket, replyPacket.ResultCode)
}
return
}
totalWritten += int64(n)
}
@ -509,6 +530,7 @@ func (rc *RemoteCacheClient) Put(ctx context.Context, reqId, key string, r io.Re
} else if err != nil {
return err
} else if totalWritten == length {
log.LogDebugf(" total written %v write size %v to fg", totalWritten, readSize)
break
}
}
@ -518,13 +540,29 @@ func (rc *RemoteCacheClient) Put(ctx context.Context, reqId, key string, r io.Re
return nil
}
func (rc *RemoteCacheClient) deleteRemoteBlock(key string, conn *net.TCPConn) {
func (rc *RemoteCacheClient) deleteRemoteBlock(key string, addr string) {
var (
err error
conn *net.TCPConn
)
if conn, err = rc.conns.GetConnect(addr); err != nil {
log.LogWarnf("FlashGroup delete: get connection to curr addr failed, addr(%v) key(%v) err(%v)", addr, key, err)
return
}
defer func() {
rc.conns.PutConnect(conn, err != nil)
}()
p := proto.NewPacketReqID()
p.Data = ([]byte)(key)
p.Opcode = proto.OpFlashNodeCacheDelete
p.Size = uint32(len(p.Data))
if err := p.WriteToConn(conn); err != nil {
log.LogWarnf("FlashGroup put: failed to write to addr(%v) err(%v)", conn.RemoteAddr().String(), err)
if err = p.WriteToConn(conn); err != nil {
log.LogWarnf("FlashGroup delete: failed to write to addr(%v) err(%v)", conn.RemoteAddr().String(), err)
return
}
replyPacket := NewFlashCacheReply()
if err = replyPacket.ReadFromConn(conn, proto.ReadDeadlineTime); err != nil {
log.LogWarnf("FlashGroup delete: failed to ReadFromConn, replyPacket(%v), fg host(%v) err(%v)", replyPacket, addr, err)
return
}
}

View File

@ -22,7 +22,6 @@ import (
"regexp"
"strings"
"github.com/cubefs/cubefs/depends/tiglabs/raft/util"
"github.com/cubefs/cubefs/util/log"
)
@ -41,7 +40,7 @@ const (
BlockCount = 1024
BlockSize = 65536 * 2
ReadBlockSize = BlockSize
RepairReadBlockSize = 512 * util.KB
RepairReadBlockSize = 512 * KB
CacheReadBlockSize = 4 * MB
PerBlockCrcSize = 4
ExtentSize = BlockCount * BlockSize
@ -57,7 +56,7 @@ const (
)
const (
PageSize = 4 * util.KB
PageSize = 4 * KB
FallocFLKeepSize = 1
FallocFLPunchHole = 2
)