diff --git a/datanode/partition_op_by_raft.go b/datanode/partition_op_by_raft.go index b27a2fc50..bdd7045cc 100644 --- a/datanode/partition_op_by_raft.go +++ b/datanode/partition_op_by_raft.go @@ -256,6 +256,9 @@ func (dp *DataPartition) ApplyRandomWrite(command []byte, raftApplyID uint64) (r log.LogErrorf("[ApplyRandomWrite] ApplyID(%v) Partition(%v) unmarshal failed(%v)", raftApplyID, dp.partitionID, err) return } + if opItem.size == 0 || opItem.data == nil { + return + } log.LogDebugf("[ApplyRandomWrite] ApplyID(%v) Partition(%v)_Extent(%v)_ExtentOffset(%v)_Size(%v)", raftApplyID, dp.partitionID, opItem.extentID, opItem.offset, opItem.size) diff --git a/datanode/server.go b/datanode/server.go index 812243ce0..7aa02f877 100644 --- a/datanode/server.go +++ b/datanode/server.go @@ -926,6 +926,7 @@ func (s *DataNode) registerHandler() { http.HandleFunc("/getRaftPeers", s.getRaftPeers) http.HandleFunc("/setGOGC", s.setGOGC) http.HandleFunc("/getGOGC", s.getGOGC) + http.HandleFunc("/triggerRaftLogRotate", s.triggerRaftLogRotate) } func (s *DataNode) startTCPService() (err error) { diff --git a/datanode/server_handler.go b/datanode/server_handler.go index b6440c18e..0b3a9b2af 100644 --- a/datanode/server_handler.go +++ b/datanode/server_handler.go @@ -868,3 +868,28 @@ func (s *DataNode) getGOGC(w http.ResponseWriter, r *http.Request) { data := fmt.Sprintf("gogc value is %v", s.gogcValue) s.buildSuccessResp(w, data) } + +func (s *DataNode) triggerRaftLogRotate(w http.ResponseWriter, r *http.Request) { + val, err := MarshalRandWriteRaftLog(proto.OpRandomWrite, 0, 0, 0, nil, 0) + if err != nil { + log.LogErrorf("action[triggerRaftLogRotate] marshal error %v", err) + s.buildFailureResp(w, http.StatusBadRequest, err.Error()) + return + } + dataPartitions := s.space.getPartitions() + trigger := 0 + for _, dp := range dataPartitions { + if !dp.raftPartition.IsRaftLeader() { + continue + } + _, err = dp.Submit(val) + if err != nil { + log.LogErrorf("action[triggerRaftLogRotate] submit error %v", err) + s.buildFailureResp(w, http.StatusBadRequest, err.Error()) + return + } + trigger += 1 + } + + s.buildSuccessResp(w, fmt.Sprintf("trigger dp(%d) raft log rotate successfully.", trigger)) +} diff --git a/depends/tiglabs/raft/storage/wal/config.go b/depends/tiglabs/raft/storage/wal/config.go index 20677bf22..d9daa57bc 100644 --- a/depends/tiglabs/raft/storage/wal/config.go +++ b/depends/tiglabs/raft/storage/wal/config.go @@ -19,8 +19,8 @@ import "github.com/cubefs/cubefs/depends/tiglabs/raft/util" const ( DefaultFileCacheCapacity = 2 DefaultFileSize = 32 * util.MB - MinFileSize = 1 * util.MB - MaxRotateInterval = 86400 + MinFileSize = 128 * util.KB + MaxRotateInterval = 3600 DefaultSync = false ) diff --git a/depends/tiglabs/raft/storage/wal/log_storage.go b/depends/tiglabs/raft/storage/wal/log_storage.go index 9c915ac03..b5b7b6ede 100644 --- a/depends/tiglabs/raft/storage/wal/log_storage.go +++ b/depends/tiglabs/raft/storage/wal/log_storage.go @@ -44,7 +44,7 @@ func openLogStorage(dir string, s *Storage) (*logEntryStorage, error) { s: s, dir: dir, filesize: s.c.GetFileSize(), - rotateTime: timeutil.GetCurrentTimeUnix(), + rotateTime: timeutil.GetCurrentTimeUnix() - MaxRotateInterval, nextFileSeq: 1, }