Refactor: add autoRepait limit on datanode

Signed-off-by: awzhgw <guowl18702995996@gmail.com>
This commit is contained in:
awzhgw 2020-08-18 19:52:24 +08:00
parent e3c98d07f6
commit af64301e02
5 changed files with 84 additions and 19 deletions

View File

@ -525,9 +525,6 @@ func (dp *DataPartition) streamRepairExtent(remoteExtentInfo *storage.ExtentInfo
remoteExtentInfo.Source, remoteExtentInfo.Size, currFixOffset, request.GetUniqueLogId(), reply.GetUniqueLogId())
return errors.Trace(err, "streamRepairExtent receive data error")
}
AutoRepairLimiterWait()
isEmptyResponse := false
// Write it to local extent file
if storage.IsTinyExtent(uint64(localExtentInfo.FileID)) {

View File

@ -2,20 +2,49 @@ package datanode
import (
"context"
"fmt"
"golang.org/x/time/rate"
)
var (
deleteLimiteRater = rate.NewLimiter(rate.Inf, defaultMarkDeleteLimitBurst)
autoRepairLimiteRater = rate.NewLimiter(rate.Inf, 512)
deleteLimiteRater = rate.NewLimiter(rate.Inf, defaultMarkDeleteLimitBurst)
MaxExtentRepairLimit = 20000
MinExtentRepairLimit = 5
extentRepairLimiteRater = make(chan struct{}, MaxExtentRepairLimit)
)
func AutoRepairLimiterWait() (err error) {
ctx := context.Background()
autoRepairLimiteRater.Wait(ctx)
func requestDoExtentRepair() (err error) {
err = fmt.Errorf("cannot do extentRepair")
select {
case <-extentRepairLimiteRater:
return nil
default:
return
}
return
}
func fininshDoExtentRepair() {
select {
case extentRepairLimiteRater <- struct{}{}:
return
default:
return
}
}
func setDoExtentRepair(value int) {
close(extentRepairLimiteRater)
if value > MaxExtentRepairLimit {
value = MaxExtentRepairLimit
}
if value < MinExtentRepairLimit {
value = MinExtentRepairLimit
}
extentRepairLimiteRater = make(chan struct{}, value)
}
func DeleteLimiterWait() {
ctx := context.Background()
deleteLimiteRater.Wait(ctx)

View File

@ -42,7 +42,7 @@ func (m *DataNode) updateNodeInfo() {
return
}
setLimiter(deleteLimiteRater, clusterInfo.DataNodeDeleteLimitRate)
setLimiter(autoRepairLimiteRater, clusterInfo.DataNodeAutoRepairLimitRate)
setDoExtentRepair(int(clusterInfo.DataNodeAutoRepairLimitRate))
log.LogInfof("updateNodeInfo from master:"+
"deleteLimite(%v),autoRepairLimit(%v)", clusterInfo.DataNodeDeleteLimitRate,
clusterInfo.DataNodeAutoRepairLimitRate)

View File

@ -714,8 +714,7 @@ func (dp *DataPartition) pushSyncDeleteRecordFromLeaderMesg() bool {
return false
}
func (dp *DataPartition)consumeTinyDeleteRecordFromLeaderMesg() {
func (dp *DataPartition) consumeTinyDeleteRecordFromLeaderMesg() {
select {
case <-dp.Disk().syncTinyDeleteRecordFromLeaderOnEveryDisk:
return
@ -730,7 +729,7 @@ func (dp *DataPartition) doStreamFixTinyDeleteRecord(repairTask *DataPartitionRe
err error
conn *net.TCPConn
)
if !dp.pushSyncDeleteRecordFromLeaderMesg(){
if !dp.pushSyncDeleteRecordFromLeaderMesg() {
return
}

View File

@ -73,11 +73,11 @@ func (s *DataNode) OperatePacket(p *repl.Packet, c *net.TCPConn) (err error) {
case proto.OpStreamRead:
s.handleStreamReadPacket(p, c, StreamRead)
case proto.OpStreamFollowerRead:
s.handleExtentRepaiReadPacket(p, c, StreamRead)
s.extentRepaiReadPacket(p, c, StreamRead)
case proto.OpExtentRepairRead:
s.handleExtentRepaiReadPacket(p, c, RepairRead)
s.handleExtentRepairReadPacket(p, c, RepairRead)
case proto.OpTinyExtentRepairRead:
s.handleTinyExtentRepairRead(p, c)
s.handleTinyExtentRepairReadPacket(p, c)
case proto.OpMarkDelete:
s.handleMarkDeletePacket(p, c)
case proto.OpBatchDeleteExtent:
@ -464,12 +464,52 @@ func (s *DataNode) handleStreamReadPacket(p *repl.Packet, connect net.Conn, isRe
if err = partition.CheckLeader(p, connect); err != nil {
return
}
s.handleExtentRepaiReadPacket(p, connect, isRepairRead)
s.extentRepaiReadPacket(p, connect, isRepairRead)
return
}
func (s *DataNode) handleExtentRepaiReadPacket(p *repl.Packet, connect net.Conn, isRepairRead bool) {
func (s *DataNode) handleExtentRepairReadPacket(p *repl.Packet, connect net.Conn, isRepairRead bool) {
var (
err error
)
defer func() {
if err != nil {
p.PackErrorBody(ActionStreamRead, err.Error())
p.WriteToConn(connect)
}
fininshDoExtentRepair()
}()
err = requestDoExtentRepair()
if err != nil {
return
}
s.extentRepaiReadPacket(p, connect, isRepairRead)
}
func (s *DataNode) handleTinyExtentRepairReadPacket(p *repl.Packet, connect net.Conn) {
var (
err error
)
defer func() {
if err != nil {
p.PackErrorBody(ActionStreamRead, err.Error())
p.WriteToConn(connect)
}
fininshDoExtentRepair()
}()
err = requestDoExtentRepair()
if err != nil {
return
}
s.tinyExtentRepairRead(p, connect)
}
func (s *DataNode) extentRepaiReadPacket(p *repl.Packet, connect net.Conn, isRepairRead bool) {
var (
err error
)
@ -577,8 +617,8 @@ func (s *DataNode) attachAvaliSizeOnTinyExtentRepairRead(reply *repl.Packet, ava
binary.BigEndian.PutUint64(reply.Arg[9:17], avaliSize)
}
// Handle handleTinyExtentRepairRead packet.
func (s *DataNode) handleTinyExtentRepairRead(request *repl.Packet, connect net.Conn) {
// Handle tinyExtentRepairRead packet.
func (s *DataNode) tinyExtentRepairRead(request *repl.Packet, connect net.Conn) {
var (
err error
needReplySize int64