mirror of
https://github.com/cubefs/cubefs.git
synced 2026-08-02 02:00:56 +00:00
feat(clustermgr): raft filewal cleans trash log by itself
with #1000543617 Signed-off-by: tangdeyi <tangdeyi@oppo.com>
This commit is contained in:
parent
8b085fbac6
commit
6499ead6d5
@ -57,6 +57,8 @@ type Config struct {
|
||||
|
||||
ProposeTimeout int `json:"propose_timeout"`
|
||||
|
||||
TrashLogReserveNum int `json:"trash_log_reserve_num"`
|
||||
|
||||
// if true, follower raft will not forward the proposal to leader.
|
||||
DisableProposalForwarding bool `json:"disable_proposal_forwarding"`
|
||||
|
||||
|
||||
@ -22,10 +22,11 @@ import (
|
||||
"sync/atomic"
|
||||
"time"
|
||||
|
||||
"github.com/cubefs/cubefs/blobstore/common/raftserver/wal"
|
||||
"github.com/cubefs/cubefs/blobstore/util/log"
|
||||
"go.etcd.io/etcd/raft/v3"
|
||||
pb "go.etcd.io/etcd/raft/v3/raftpb"
|
||||
|
||||
"github.com/cubefs/cubefs/blobstore/common/raftserver/wal"
|
||||
"github.com/cubefs/cubefs/blobstore/util/log"
|
||||
)
|
||||
|
||||
const (
|
||||
@ -114,7 +115,7 @@ func NewRaftServer(cfg *Config) (RaftServer, error) {
|
||||
rs.readNotifier.Store(newReadIndexNotifier())
|
||||
|
||||
begin := time.Now()
|
||||
store, err := NewRaftStorage(cfg.WalDir, cfg.WalSync, cfg.UseRocksdb, cfg.NodeId, rs.sm, rs.shotter)
|
||||
store, err := NewRaftStorage(cfg.WalDir, cfg.WalSync, cfg.UseRocksdb, cfg.NodeId, cfg.TrashLogReserveNum, rs.sm, rs.shotter)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
@ -39,7 +39,7 @@ type raftStorage struct {
|
||||
snapIndex uint64
|
||||
}
|
||||
|
||||
func NewRaftStorage(walDir string, sync bool, use_rocksdb bool, nodeId uint64, sm StateMachine, shotter *snapshotter) (*raftStorage, error) {
|
||||
func NewRaftStorage(walDir string, sync bool, use_rocksdb bool, nodeId uint64, trashLogReserveNum int, sm StateMachine, shotter *snapshotter) (*raftStorage, error) {
|
||||
rs := &raftStorage{
|
||||
nodeId: nodeId,
|
||||
shotter: shotter,
|
||||
@ -55,7 +55,7 @@ func NewRaftStorage(walDir string, sync bool, use_rocksdb bool, nodeId uint64, s
|
||||
if use_rocksdb {
|
||||
w, err = wal.OpenRocksdbWal(walDir)
|
||||
} else {
|
||||
w, err = wal.OpenWal(walDir, sync)
|
||||
w, err = wal.OpenWal(walDir, sync, trashLogReserveNum)
|
||||
}
|
||||
if err != nil {
|
||||
return nil, err
|
||||
|
||||
@ -23,8 +23,9 @@ import (
|
||||
|
||||
"github.com/stretchr/testify/require"
|
||||
|
||||
"github.com/cubefs/cubefs/blobstore/common/raftserver/wal"
|
||||
pb "go.etcd.io/etcd/raft/v3/raftpb"
|
||||
|
||||
"github.com/cubefs/cubefs/blobstore/common/raftserver/wal"
|
||||
)
|
||||
|
||||
const (
|
||||
@ -86,7 +87,7 @@ func (sm *storeSM) LeaderChange(leader uint64, host string) {
|
||||
func TestStorage(t *testing.T) {
|
||||
{
|
||||
os.RemoveAll(walDir)
|
||||
store, err := NewRaftStorage(walDir, true, true, nodeId, &storeSM{}, newSnapshotter(5, time.Second*10))
|
||||
store, err := NewRaftStorage(walDir, true, true, nodeId, 10, &storeSM{}, newSnapshotter(5, time.Second*10))
|
||||
require.Nil(t, err)
|
||||
hs, cs, _ := store.InitialState()
|
||||
require.Equal(t, hs, pb.HardState{})
|
||||
@ -97,7 +98,7 @@ func TestStorage(t *testing.T) {
|
||||
|
||||
{
|
||||
os.RemoveAll(walDir)
|
||||
store, err := NewRaftStorage(walDir, true, true, nodeId, &storeSM{}, newSnapshotter(5, time.Second*10))
|
||||
store, err := NewRaftStorage(walDir, true, true, nodeId, 10, &storeSM{}, newSnapshotter(5, time.Second*10))
|
||||
require.Nil(t, err)
|
||||
var entries []pb.Entry
|
||||
for i := 0; i < 1000; i++ {
|
||||
|
||||
@ -26,8 +26,8 @@ const (
|
||||
TrashPath = ".trash"
|
||||
)
|
||||
|
||||
func InitPath(dir string, createTrush bool) error {
|
||||
if createTrush {
|
||||
func InitPath(dir string, createTrash bool) error {
|
||||
if createTrash {
|
||||
dir = path.Join(dir, TrashPath)
|
||||
}
|
||||
info, err := os.Stat(dir)
|
||||
|
||||
@ -15,14 +15,18 @@
|
||||
package wal
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"math"
|
||||
"os"
|
||||
"path"
|
||||
"sort"
|
||||
"syscall"
|
||||
"time"
|
||||
|
||||
pb "go.etcd.io/etcd/raft/v3/raftpb"
|
||||
|
||||
"github.com/cubefs/cubefs/blobstore/common/trace"
|
||||
)
|
||||
|
||||
const (
|
||||
@ -30,6 +34,8 @@ const (
|
||||
logfileCacheNum = 4
|
||||
)
|
||||
|
||||
var trashCleanIntervalSec = 300
|
||||
|
||||
func (w *fileWal) reload(firstIndex uint64) error {
|
||||
names, err := listLogFiles(w.dir)
|
||||
if err != nil {
|
||||
@ -324,6 +330,34 @@ func (w *fileWal) saveEntry(ent *pb.Entry) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
func (w *fileWal) run() {
|
||||
if w.trashLogReserveNum <= 0 {
|
||||
return
|
||||
}
|
||||
|
||||
ticker := time.NewTicker(time.Duration(trashCleanIntervalSec) * time.Second)
|
||||
defer ticker.Stop()
|
||||
for {
|
||||
select {
|
||||
case <-ticker.C:
|
||||
span, _ := trace.StartSpanFromContext(context.Background(), "")
|
||||
trashFiles, err := listLogFiles(path.Join(w.dir, TrashPath))
|
||||
if err != nil {
|
||||
span.Warnf("fileWal trashClean listLogFiles err %v", err)
|
||||
continue
|
||||
}
|
||||
if len(trashFiles) > w.trashLogReserveNum {
|
||||
for i := 0; i < len(trashFiles)-w.trashLogReserveNum; i++ {
|
||||
err = os.Remove(path.Join(path.Join(w.dir, TrashPath), trashFiles[i].String()))
|
||||
span.Warnf("fileWal trashClean remove file %s, err %v", trashFiles[i].String(), err)
|
||||
}
|
||||
}
|
||||
case <-w.closeCh:
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func isErrInvalid(err error) bool {
|
||||
if err == os.ErrInvalid {
|
||||
return true
|
||||
|
||||
@ -20,6 +20,7 @@ import (
|
||||
"runtime"
|
||||
"syscall"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/stretchr/testify/require"
|
||||
"go.etcd.io/etcd/raft/v3/raftpb"
|
||||
@ -30,6 +31,7 @@ func openLogStorage(dir string) *fileWal {
|
||||
return &fileWal{
|
||||
dir: dir,
|
||||
nextFileSeq: 1,
|
||||
closeCh: make(chan struct{}),
|
||||
cache: newLogFileCache(logfileCacheNum,
|
||||
func(name logName) (*logFile, error) {
|
||||
lf, err := openLogFile(dir, name, false)
|
||||
@ -195,3 +197,51 @@ func TestLogStoragetruncateBack(t *testing.T) {
|
||||
require.NotNil(t, err)
|
||||
ls.Close()
|
||||
}
|
||||
|
||||
func TestRun(t *testing.T) {
|
||||
dir := "/tmp/raftwal-run-" + string(genRandomBytes(8))
|
||||
clear := func() {
|
||||
os.RemoveAll(dir)
|
||||
}
|
||||
clear()
|
||||
defer clear()
|
||||
|
||||
err := InitPath(dir, true)
|
||||
require.Nil(t, err)
|
||||
|
||||
trashDir := dir + "/" + TrashPath
|
||||
trashLogReserveNum := 3
|
||||
|
||||
for i := 0; i < 5; i++ {
|
||||
name := logName{sequence: uint64(i), index: uint64(i * 100)}
|
||||
f, createErr := os.Create(trashDir + "/" + name.String())
|
||||
require.Nil(t, createErr)
|
||||
f.Close()
|
||||
}
|
||||
|
||||
files, err := listLogFiles(trashDir)
|
||||
require.Nil(t, err)
|
||||
require.Equal(t, 5, len(files))
|
||||
|
||||
trashCleanIntervalSec = 1
|
||||
|
||||
w := &fileWal{
|
||||
dir: dir,
|
||||
closeCh: make(chan struct{}),
|
||||
trashLogReserveNum: trashLogReserveNum,
|
||||
}
|
||||
|
||||
go w.run()
|
||||
|
||||
time.Sleep(5 * time.Second)
|
||||
|
||||
close(w.closeCh)
|
||||
|
||||
files, err = listLogFiles(trashDir)
|
||||
require.Nil(t, err)
|
||||
require.Equal(t, trashLogReserveNum, len(files))
|
||||
|
||||
require.Equal(t, uint64(2), files[0].sequence)
|
||||
require.Equal(t, uint64(3), files[1].sequence)
|
||||
require.Equal(t, uint64(4), files[2].sequence)
|
||||
}
|
||||
|
||||
@ -19,9 +19,10 @@ import (
|
||||
"runtime"
|
||||
"sync"
|
||||
|
||||
"github.com/cubefs/cubefs/blobstore/util/log"
|
||||
"go.etcd.io/etcd/raft/v3"
|
||||
pb "go.etcd.io/etcd/raft/v3/raftpb"
|
||||
|
||||
"github.com/cubefs/cubefs/blobstore/util/log"
|
||||
)
|
||||
|
||||
type Snapshot struct {
|
||||
@ -54,10 +55,13 @@ type fileWal struct {
|
||||
st Snapshot
|
||||
mt *meta
|
||||
once sync.Once
|
||||
closeCh chan struct{}
|
||||
|
||||
trashLogReserveNum int
|
||||
}
|
||||
|
||||
// OpenWal
|
||||
func OpenWal(dir string, sync bool) (Wal, error) {
|
||||
func OpenWal(dir string, sync bool, trashLogReserveNum int) (Wal, error) {
|
||||
dir = path.Clean(dir)
|
||||
if err := InitPath(dir, true); err != nil {
|
||||
return nil, err
|
||||
@ -84,6 +88,9 @@ func OpenWal(dir string, sync bool) (Wal, error) {
|
||||
hs: hs,
|
||||
st: st,
|
||||
mt: mt,
|
||||
|
||||
closeCh: make(chan struct{}),
|
||||
trashLogReserveNum: trashLogReserveNum,
|
||||
}
|
||||
|
||||
err = w.reload(st.Index + 1)
|
||||
@ -96,6 +103,8 @@ func OpenWal(dir string, sync bool) (Wal, error) {
|
||||
w.hs.Commit = w.LastIndex()
|
||||
}
|
||||
|
||||
go w.run()
|
||||
|
||||
return w, nil
|
||||
}
|
||||
|
||||
@ -200,5 +209,6 @@ func (w *fileWal) Close() {
|
||||
w.cache.Close()
|
||||
w.last.Close()
|
||||
w.mt.Close()
|
||||
close(w.closeCh)
|
||||
})
|
||||
}
|
||||
|
||||
@ -54,7 +54,7 @@ func TestRaftWal(t *testing.T) {
|
||||
}
|
||||
clear()
|
||||
defer clear()
|
||||
wal, err := OpenWal(dir, true)
|
||||
wal, err := OpenWal(dir, true, 10)
|
||||
require.Nil(t, err)
|
||||
|
||||
entries := make([]pb.Entry, 0, 100000)
|
||||
@ -159,7 +159,7 @@ func TestRaftWal(t *testing.T) {
|
||||
require.Nil(t, err)
|
||||
wal.Close()
|
||||
|
||||
wal, err = OpenWal(dir, true)
|
||||
wal, err = OpenWal(dir, true, 10)
|
||||
require.Nil(t, err)
|
||||
wal.Close()
|
||||
}
|
||||
|
||||
Loading…
Reference in New Issue
Block a user