mirror of
https://github.com/cubefs/cubefs.git
synced 2026-08-02 02:00:56 +00:00
enhance(client): reconstruct process of try init extentHandler by last ek
Signed-off-by: leonrayang <chl696@sina.com>
This commit is contained in:
parent
88a5b8247d
commit
9f1bbe8e67
@ -83,7 +83,7 @@ func (s *DataNode) OperatePacket(p *repl.Packet, c net.Conn) (err error) {
|
||||
tpLabels map[string]string
|
||||
tpObject *exporter.TimePointCount
|
||||
)
|
||||
|
||||
log.LogDebugf("action[OperatePacket] %v, pack [%v]", p.GetOpMsg(), p)
|
||||
shallDegrade := p.ShallDegrade()
|
||||
sz := p.Size
|
||||
if !shallDegrade {
|
||||
|
||||
@ -354,8 +354,8 @@ func (p *Packet) GetCopy() *Packet {
|
||||
}
|
||||
|
||||
func (p *Packet) String() string {
|
||||
return fmt.Sprintf("ReqID(%v)Op(%v)PartitionID(%v)ResultCode(%v)ExID(%v)ExtOffset(%v)KernelOff(%v)Type(%v)Seq(%v)",
|
||||
p.ReqID, p.GetOpMsg(), p.PartitionID, p.GetResultMsg(), p.ExtentID, p.ExtentOffset, p.KernelOffset, p.ExtentType, p.VerSeq)
|
||||
return fmt.Sprintf("ReqID(%v)Op(%v)PartitionID(%v)ResultCode(%v)ExID(%v)ExtOffset(%v)KernelOff(%v)Type(%v)Seq(%v)Size(%v)",
|
||||
p.ReqID, p.GetOpMsg(), p.PartitionID, p.GetResultMsg(), p.ExtentID, p.ExtentOffset, p.KernelOffset, p.ExtentType, p.VerSeq, p.Size)
|
||||
}
|
||||
|
||||
// GetStoreType returns the store type.
|
||||
|
||||
@ -506,7 +506,7 @@ func (cache *ExtentCache) PrepareWriteRequests(offset, size int, data []byte) []
|
||||
cache.root.DescendLessOrEqual(pivot, func(i btree.Item) bool {
|
||||
ek := i.(*proto.ExtentKey)
|
||||
lower.FileOffset = ek.FileOffset
|
||||
// log.LogDebugf("action[ExtentCache.PrepareWriteRequests] ek [%v], pivot[%v]", ek, pivot)
|
||||
log.LogDebugf("action[ExtentCache.PrepareWriteRequests] ek [%v], pivot[%v]", ek, pivot)
|
||||
return false
|
||||
})
|
||||
|
||||
@ -515,7 +515,7 @@ func (cache *ExtentCache) PrepareWriteRequests(offset, size int, data []byte) []
|
||||
ekStart := int(ek.FileOffset)
|
||||
ekEnd := int(ek.FileOffset) + int(ek.Size)
|
||||
|
||||
// log.LogDebugf("PrepareWriteRequests: ino(%v) start(%v) end(%v) ekStart(%v) ekEnd(%v)", cache.inode, start, end, ekStart, ekEnd)
|
||||
log.LogDebugf("action[ExtentCache.PrepareWriteRequests]: ino(%v) start(%v) end(%v) ekStart(%v) ekEnd(%v)", cache.inode, start, end, ekStart, ekEnd)
|
||||
|
||||
if start <= ekStart {
|
||||
if end <= ekStart {
|
||||
@ -554,7 +554,7 @@ func (cache *ExtentCache) PrepareWriteRequests(offset, size int, data []byte) []
|
||||
}
|
||||
})
|
||||
|
||||
// log.LogDebugf("PrepareWriteRequests: ino(%v) start(%v) end(%v)", cache.inode, start, end)
|
||||
log.LogDebugf("PrepareWriteRequests: ino(%v) start(%v) end(%v)", cache.inode, start, end)
|
||||
if start < end {
|
||||
// add hole (start, end)
|
||||
req := NewExtentRequest(start, end-start, data[start-offset:end-offset], nil)
|
||||
|
||||
@ -113,7 +113,7 @@ type ExtentHandler struct {
|
||||
|
||||
// NewExtentHandler returns a new extent handler.
|
||||
func NewExtentHandler(stream *Streamer, offset int, storeMode int, size int) *ExtentHandler {
|
||||
// log.LogDebugf("NewExtentHandler stack(%v)", string(debug.Stack()))
|
||||
// log.LogDebugf("NewExtentHandler stack(%v)", string(debug.Stack()))
|
||||
eh := &ExtentHandler{
|
||||
stream: stream,
|
||||
id: GetExtentHandlerID(),
|
||||
@ -137,8 +137,8 @@ func NewExtentHandler(stream *Streamer, offset int, storeMode int, size int) *Ex
|
||||
|
||||
// String returns the string format of the extent handler.
|
||||
func (eh *ExtentHandler) String() string {
|
||||
return fmt.Sprintf("ExtentHandler{ID(%v)Inode(%v)FileOffset(%v)StoreMode(%v)Status(%v)}",
|
||||
eh.id, eh.inode, eh.fileOffset, eh.storeMode, eh.status)
|
||||
return fmt.Sprintf("ExtentHandler{ID(%v)Inode(%v)FileOffset(%v)StoreMode(%v)Status(%v)Dp(%v)}",
|
||||
eh.id, eh.inode, eh.fileOffset, eh.storeMode, eh.status, eh.dp)
|
||||
}
|
||||
|
||||
func (eh *ExtentHandler) write(data []byte, offset, size int, direct bool) (ek *proto.ExtentKey, err error) {
|
||||
@ -216,7 +216,7 @@ func (eh *ExtentHandler) sender() {
|
||||
// case <-t.C:
|
||||
// log.LogDebugf("sender alive: eh(%v) inflight(%v)", eh, atomic.LoadInt32(&eh.inflight))
|
||||
case packet := <-eh.request:
|
||||
//log.LogDebugf("ExtentHandler sender begin: eh(%v) packet(%v)", eh, packet.GetUniqueLogId())
|
||||
log.LogDebugf("ExtentHandler sender begin: eh(%v) packet(%v)", eh, packet)
|
||||
if eh.getStatus() >= ExtentStatusRecovery {
|
||||
log.LogWarnf("sender in recovery: eh(%v) packet(%v)", eh, packet)
|
||||
eh.reply <- packet
|
||||
@ -259,7 +259,7 @@ func (eh *ExtentHandler) sender() {
|
||||
}
|
||||
packet.StartT = time.Now().UnixNano()
|
||||
|
||||
//log.LogDebugf("ExtentHandler sender: extent allocated, eh(%v) dp(%v) extID(%v) packet(%v)", eh, eh.dp, eh.extID, packet.GetUniqueLogId())
|
||||
log.LogDebugf("ExtentHandler sender: extent allocated, eh(%v) dp(%v) extID(%v) packet(%v)", eh, eh.dp, eh.extID, packet.GetUniqueLogId())
|
||||
|
||||
if err = packet.writeToConn(eh.conn); err != nil {
|
||||
log.LogWarnf("sender writeTo: failed, eh(%v) err(%v) packet(%v)", eh, err, packet)
|
||||
@ -284,9 +284,9 @@ func (eh *ExtentHandler) receiver() {
|
||||
// case <-t.C:
|
||||
// log.LogDebugf("receiver alive: eh(%v) inflight(%v)", eh, atomic.LoadInt32(&eh.inflight))
|
||||
case packet := <-eh.reply:
|
||||
//log.LogDebugf("receiver begin: eh(%v) packet(%v)", eh, packet.GetUniqueLogId())
|
||||
log.LogDebugf("receiver begin: eh(%v) packet(%v)", eh, packet.GetUniqueLogId())
|
||||
eh.processReply(packet)
|
||||
//log.LogDebugf("receiver end: eh(%v) packet(%v)", eh, packet.GetUniqueLogId())
|
||||
log.LogDebugf("receiver end: eh(%v) packet(%v)", eh, packet.GetUniqueLogId())
|
||||
case <-eh.doneReceiver:
|
||||
log.LogDebugf("receiver done: eh(%v) size(%v) ek(%v)", eh, eh.size, eh.key)
|
||||
return
|
||||
@ -325,6 +325,7 @@ func (eh *ExtentHandler) processReply(packet *Packet) {
|
||||
if reply.ResultCode != proto.OpOk {
|
||||
if reply.ResultCode != proto.ErrCodeVersionOpError {
|
||||
errmsg := fmt.Sprintf("reply NOK: reply(%v)", reply)
|
||||
log.LogDebugf("processReply packet (%v) errmsg (%v)", packet, errmsg)
|
||||
eh.processReplyError(packet, errmsg)
|
||||
return
|
||||
}
|
||||
@ -367,7 +368,7 @@ func (eh *ExtentHandler) processReply(packet *Packet) {
|
||||
ExtentOffset: extOffset,
|
||||
Size: packet.Size,
|
||||
SnapInfo: &proto.ExtSnapInfo{
|
||||
VerSeq: eh.stream.verSeq,
|
||||
VerSeq: reply.VerSeq,
|
||||
},
|
||||
}
|
||||
} else {
|
||||
@ -654,12 +655,12 @@ func (eh *ExtentHandler) setClosed() bool {
|
||||
}
|
||||
|
||||
func (eh *ExtentHandler) setRecovery() bool {
|
||||
// log.LogDebugf("action[ExtentHandler.setRecovery] stack (%v)", string(debug.Stack()))
|
||||
// log.LogDebugf("action[ExtentHandler.setRecovery] stack (%v)", string(debug.Stack()))
|
||||
return atomic.CompareAndSwapInt32(&eh.status, ExtentStatusClosed, ExtentStatusRecovery)
|
||||
}
|
||||
|
||||
func (eh *ExtentHandler) setError() bool {
|
||||
// log.LogDebugf("action[ExtentHandler.setError] stack (%v)", string(debug.Stack()))
|
||||
// log.LogDebugf("action[ExtentHandler.setError] stack (%v)", string(debug.Stack()))
|
||||
if proto.IsHot(eh.stream.client.volumeType) {
|
||||
atomic.StoreInt32(&eh.stream.status, StreamerError)
|
||||
}
|
||||
|
||||
@ -693,54 +693,80 @@ func (s *Streamer) doOverwrite(req *ExtentRequest, direct bool) (total int, err
|
||||
|
||||
func (s *Streamer) tryInitExtentHandlerByLastEk(offset, size int) (isLastEkVerNotEqual bool) {
|
||||
storeMode := s.GetStoreMod(offset, size)
|
||||
getEndEkFunc := func() *proto.ExtentKey {
|
||||
if ek := s.extents.GetEndForAppendWrite(uint64(offset), s.verSeq, false); ek != nil && !storage.IsTinyExtent(ek.ExtentId) {
|
||||
return ek
|
||||
}
|
||||
return nil
|
||||
}
|
||||
initExtentHandlerFunc := func(currentEK *proto.ExtentKey) {
|
||||
if currentEK.GetSeq() != s.verSeq {
|
||||
log.LogDebugf("tryInitExtentHandlerByLastEk. exist ek seq %v vs request seq %v", currentEK.GetSeq(), s.verSeq)
|
||||
if int(currentEK.ExtentOffset)+int(currentEK.Size)+size > util.ExtentSize {
|
||||
s.closeOpenHandler()
|
||||
return
|
||||
}
|
||||
isLastEkVerNotEqual = true
|
||||
}
|
||||
|
||||
// && (s.handler == nil || s.handler != nil && s.handler.fileOffset+s.handler.size != offset) delete ??
|
||||
if storeMode == proto.NormalExtentType && (s.handler == nil || s.handler != nil && s.handler.fileOffset+s.handler.size != offset) {
|
||||
if currentEK := s.extents.GetEndForAppendWrite(uint64(offset), s.verSeq, false); currentEK != nil && !storage.IsTinyExtent(currentEK.ExtentId) {
|
||||
if currentEK.GetSeq() != s.verSeq {
|
||||
log.LogDebugf("tryInitExtentHandlerByLastEk. exist ek seq %v vs request seq %v", currentEK.GetSeq(), s.verSeq)
|
||||
if int(currentEK.ExtentOffset)+int(currentEK.Size)+size > util.ExtentSize {
|
||||
s.closeOpenHandler()
|
||||
log.LogDebugf("tryInitExtentHandlerByLastEk: found ek in ExtentCache, extent_id(%v) req_offset(%v) req_size(%v), currentEK [%v] streamer seq %v",
|
||||
currentEK.ExtentId, offset, size, currentEK, s.verSeq)
|
||||
_, pidErr := s.client.dataWrapper.GetDataPartition(currentEK.PartitionId)
|
||||
if pidErr == nil {
|
||||
seq := currentEK.GetSeq()
|
||||
if isLastEkVerNotEqual {
|
||||
seq = s.verSeq
|
||||
}
|
||||
handler := NewExtentHandler(s, int(currentEK.FileOffset), storeMode, int(currentEK.Size))
|
||||
handler.key = &proto.ExtentKey{
|
||||
FileOffset: currentEK.FileOffset,
|
||||
PartitionId: currentEK.PartitionId,
|
||||
ExtentId: currentEK.ExtentId,
|
||||
ExtentOffset: currentEK.ExtentOffset,
|
||||
Size: currentEK.Size,
|
||||
SnapInfo: &proto.ExtSnapInfo{
|
||||
VerSeq: seq,
|
||||
},
|
||||
}
|
||||
//handler.dp = dp
|
||||
|
||||
if s.handler != nil {
|
||||
log.LogDebugf("tryInitExtentHandlerByLastEk: close old handler, currentEK.PartitionId(%v)",
|
||||
currentEK.PartitionId)
|
||||
s.closeOpenHandler()
|
||||
}
|
||||
|
||||
s.handler = handler
|
||||
s.dirty = false
|
||||
log.LogDebugf("tryInitExtentHandlerByLastEk: currentEK.PartitionId(%v) found", currentEK.PartitionId)
|
||||
} else {
|
||||
log.LogDebugf("tryInitExtentHandlerByLastEk: currentEK.PartitionId(%v) not found", currentEK.PartitionId)
|
||||
}
|
||||
}
|
||||
|
||||
if storeMode == proto.NormalExtentType {
|
||||
if s.handler == nil {
|
||||
log.LogDebugf("tryInitExtentHandlerByLastEk: handler nil")
|
||||
if ek := getEndEkFunc(); ek != nil {
|
||||
initExtentHandlerFunc(ek)
|
||||
}
|
||||
} else {
|
||||
if s.handler.fileOffset+s.handler.size == offset {
|
||||
if s.extents.Max().GetSeq() == s.verSeq {
|
||||
log.LogDebugf("tryInitExtentHandlerByLastEk: seq %vequal", s.verSeq)
|
||||
return
|
||||
}
|
||||
isLastEkVerNotEqual = true
|
||||
}
|
||||
|
||||
log.LogDebugf("tryInitExtentHandlerByLastEk: found ek in ExtentCache, extent_id(%v) offset(%v) size(%v), ekoffset(%v) eksize(%v) exist ek seq %v vs request seq %v",
|
||||
currentEK.ExtentId, offset, size, currentEK.FileOffset, currentEK.Size, currentEK.GetSeq(), s.verSeq)
|
||||
_, pidErr := s.client.dataWrapper.GetDataPartition(currentEK.PartitionId)
|
||||
if pidErr == nil {
|
||||
seq := currentEK.GetSeq()
|
||||
if isLastEkVerNotEqual {
|
||||
seq = s.verSeq
|
||||
}
|
||||
handler := NewExtentHandler(s, int(currentEK.FileOffset), storeMode, int(currentEK.Size))
|
||||
handler.key = &proto.ExtentKey{
|
||||
FileOffset: currentEK.FileOffset,
|
||||
PartitionId: currentEK.PartitionId,
|
||||
ExtentId: currentEK.ExtentId,
|
||||
ExtentOffset: currentEK.ExtentOffset,
|
||||
Size: currentEK.Size,
|
||||
SnapInfo: &proto.ExtSnapInfo{
|
||||
VerSeq: seq,
|
||||
},
|
||||
}
|
||||
|
||||
if s.handler != nil {
|
||||
log.LogDebugf("tryInitExtentHandlerByLastEk: close old handler, currentEK.PartitionId(%v)",
|
||||
currentEK.PartitionId)
|
||||
s.closeOpenHandler()
|
||||
}
|
||||
|
||||
s.handler = handler
|
||||
s.dirty = false
|
||||
log.LogDebugf("tryInitExtentHandlerByLastEk: currentEK.PartitionId(%v) found", currentEK.PartitionId)
|
||||
log.LogDebugf("tryInitExtentHandlerByLastEk: seq not equal %v:%v", s.extents.Max().GetSeq(), s.verSeq)
|
||||
initExtentHandlerFunc(s.extents.Max())
|
||||
return
|
||||
} else {
|
||||
log.LogDebugf("tryInitExtentHandlerByLastEk: currentEK.PartitionId(%v) not found", currentEK.PartitionId)
|
||||
if ek := getEndEkFunc(); ek != nil {
|
||||
log.LogDebugf("tryInitExtentHandlerByLastEk: getEndEkFunc get ek %v", ek)
|
||||
initExtentHandlerFunc(ek)
|
||||
} else {
|
||||
log.LogDebugf("tryInitExtentHandlerByLastEk: not found ek")
|
||||
}
|
||||
}
|
||||
|
||||
} else {
|
||||
log.LogDebugf("tryInitExtentHandlerByLastEk: not found ek in ExtentCache, offset(%v) size(%v)", offset, size)
|
||||
}
|
||||
}
|
||||
|
||||
@ -761,11 +787,14 @@ func (s *Streamer) doAppendWrite(data []byte, offset, size int, direct bool, reU
|
||||
if proto.IsHot(s.client.volumeType) {
|
||||
if reUseEk {
|
||||
if isLastEkVerNotEqual := s.tryInitExtentHandlerByLastEk(offset, size); isLastEkVerNotEqual {
|
||||
log.LogDebugf("doAppendWrite enter: ino(%v) tryInitExtentHandlerByLastEk worked", s.inode)
|
||||
log.LogDebugf("doAppendWrite enter: ino(%v) tryInitExtentHandlerByLastEk worked but seq not equal", s.inode)
|
||||
status = LastEKVersionNotEqual
|
||||
return
|
||||
}
|
||||
} else if s.handler != nil {
|
||||
s.closeOpenHandler()
|
||||
}
|
||||
|
||||
for i := 0; i < MaxNewHandlerRetry; i++ {
|
||||
if s.handler == nil {
|
||||
s.handler = NewExtentHandler(s, offset, storeMode, 0)
|
||||
@ -776,6 +805,8 @@ func (s *Streamer) doAppendWrite(data []byte, offset, size int, direct bool, reU
|
||||
continue
|
||||
}
|
||||
ek, err = s.handler.write(data, offset, size, direct)
|
||||
ek.SetSeq(s.verSeq)
|
||||
|
||||
if err == nil && ek != nil {
|
||||
if !s.dirty {
|
||||
s.dirtylist.Put(s.handler)
|
||||
|
||||
Loading…
Reference in New Issue
Block a user