From 81a64c5902a4f000982429cca52192e5ce64c6f4 Mon Sep 17 00:00:00 2001 From: Ben Johnson Date: Wed, 11 Nov 2020 11:17:07 -0700 Subject: [PATCH] Add RBF halting; remove time based checkpoint --- rbf/cfg/cfg.go | 30 +++++++++++------------ rbf/cfg/os.go | 2 +- rbf/db.go | 40 +++++++++++++++---------------- rbf/db_test.go | 65 +++++++++++++++++++++++++++++++++++++++++++++++++- 4 files changed, 100 insertions(+), 37 deletions(-) diff --git a/rbf/cfg/cfg.go b/rbf/cfg/cfg.go index 537158eb1..c91ad8e05 100644 --- a/rbf/cfg/cfg.go +++ b/rbf/cfg/cfg.go @@ -15,11 +15,14 @@ package cfg import ( - "time" - "github.com/spf13/pflag" ) +const ( + DefaultMinWALCheckpointSize = 1 * (1 << 20) // 1MB + DefaultMaxWALCheckpointSize = DefaultMaxWALSize / 2 +) + // Config defines externally configurable rbf options. // The separate package avoids circular import. type Config struct { @@ -30,38 +33,35 @@ type Config struct { // The maximum allowed WAL size. Required by mmap. MaxWALSize int64 + // The minimum WAL size before the WAL is copied to the DB. + MinWALCheckpointSize int64 + + // The maximum WAL size before transactions are halted to allow a checkpoint. + MaxWALCheckpointSize int64 + // Set before calling db.Open() FsyncEnabled bool // for mmap correctness testing. DoAllocZero bool - - // CheckpointEveryDur if zero means checkpoint after every write. - // Otherwise, wait and checkpoint at the next write that happens - // after CheckpointEveryDur since the previous. - CheckpointEveryDur time.Duration - - // Maximum size of a WAL write cache. - MaxWALWriteCacheSize int } func NewDefaultConfig() *Config { return &Config{ MaxSize: DefaultMaxSize, MaxWALSize: DefaultMaxWALSize, + MinWALCheckpointSize: DefaultMinWALCheckpointSize, + MaxWALCheckpointSize: DefaultMaxWALCheckpointSize, FsyncEnabled: true, - CheckpointEveryDur: time.Millisecond, - MaxWALWriteCacheSize: 1 << 20, } } func (cfg *Config) DefineFlags(flags *pflag.FlagSet) { default0 := NewDefaultConfig() - flags.IntVar(&cfg.MaxWALWriteCacheSize, "rbf-max-write-cache-size", default0.MaxWALWriteCacheSize, "RBF write cache size, in bytes") - flags.DurationVar(&cfg.CheckpointEveryDur, "rbf-checkpoint-dur", default0.CheckpointEveryDur, - "RBF checkpoint on the next write that occurs this long or more after the previous write. 0 means checkpoint after every write.") flags.Int64Var(&cfg.MaxSize, "rbf-max-db-size", default0.MaxSize, "RBF maximum size in bytes of a database file (distinct from a WAL file)") flags.Int64Var(&cfg.MaxWALSize, "rbf-max-wal-size", default0.MaxWALSize, "RBF maximum size in bytes of a WAL file (distinct from a DB file)") + flags.Int64Var(&cfg.MinWALCheckpointSize, "rbf-min-wal-checkpoint-size", default0.MinWALCheckpointSize, "RBF minimum size in bytes of a WAL file before attempting checkpoint") + flags.Int64Var(&cfg.MaxWALCheckpointSize, "rbf-max-wal-checkpoint-size", default0.MaxWALCheckpointSize, "RBF maximum size in bytes of a WAL file before forcing checkpoint") // renamed from --rbf-fsync to just --fsync because now it applies to all Tx backends. flags.BoolVar(&cfg.FsyncEnabled, "fsync", default0.FsyncEnabled, "enable fsync fully safe flush-to-disk") diff --git a/rbf/cfg/os.go b/rbf/cfg/os.go index 1dd9b1c96..28605af8b 100644 --- a/rbf/cfg/os.go +++ b/rbf/cfg/os.go @@ -24,4 +24,4 @@ const DefaultMaxSize = 4 * (1 << 30) // DefaultMaxWALSize is the default mmap size and therefore the maximum allowed // size of the WAL. The size can be increased by updating the DB.MaxWALSize // and reopening the database. This setting mainly affects virtual space usage. -const DefaultMaxWALSize = 2 * (1 << 30) +const DefaultMaxWALSize = 4 * (1 << 30) diff --git a/rbf/db.go b/rbf/db.go index c1bb4e020..a570b6f00 100644 --- a/rbf/db.go +++ b/rbf/db.go @@ -22,7 +22,6 @@ import ( "path/filepath" "sync" "syscall" - "time" "github.com/benbjohnson/immutable" "github.com/pilosa/pilosa/v2/syswrap" @@ -49,15 +48,13 @@ type DB struct { wal []byte // wal mmap walFile *os.File // wal file descriptor walPageN int // wal page count - wcache []byte // wal write cache - mu sync.RWMutex // general mutex - rwmu sync.Mutex // mutex for restricting single writer + mu sync.RWMutex // general mutex + rwmu sync.Mutex // mutex for restricting single writer + haltCond *sync.Cond // condition for resuming txs after checkpoint // Path represents the path to the database file. Path string - - lastCheckpoint time.Time } // NewDB returns a new instance of DB. @@ -70,9 +67,9 @@ func NewDB(path string, cfg *rbfcfg.Config) *DB { cfg: *cfg, txs: make(map[*Tx]struct{}), pageMap: immutable.NewMap(&uint32Hasher{}), - wcache: make([]byte, 0, cfg.MaxWALWriteCacheSize), Path: path, } + db.haltCond = sync.NewCond(&db.mu) return db } @@ -240,6 +237,9 @@ func (db *DB) checkpoint() error { db.walPageN = 0 db.pageMap = immutable.NewMap(&uint32Hasher{}) + // Notify halted tranactions that the WAL has been checkpointed. + db.haltCond.Broadcast() + return nil } @@ -446,6 +446,11 @@ func (db *DB) Begin(writable bool) (_ *Tx, err error) { return nil, ErrClosed } + // Wait for WAL size to be below threshold. + for int64(db.walPageN*PageSize) > db.cfg.MaxWALCheckpointSize { + db.haltCond.Wait() + } + tx := &Tx{ db: db, rootRecords: db.rootRecords, @@ -495,22 +500,17 @@ func (db *DB) removeTx(tx *Tx) error { delete(tx.db.txs, tx) - // Write pages from WAL to DB. - // TODO(bbj): Move this to an async goroutine. - // TODO(jea): Make the time-based checkpointing work at all, and update the - // comment in cfg/cfg.go for CheckpointEveryDur. - if tx.writable { - if db.cfg.CheckpointEveryDur == 0 || time.Since(db.lastCheckpoint) > db.cfg.CheckpointEveryDur { - if err := db.checkpoint(); err != nil { - return fmt.Errorf("checkpoint: %w", err) - } - db.lastCheckpoint = time.Now() - } - } - // Disassociate from db. tx.db = nil + // Write pages from WAL to DB. + // TODO(bbj): Move this to an async goroutine. + if len(db.txs) == 0 && db.walSize() > db.cfg.MinWALCheckpointSize { + if err := db.checkpoint(); err != nil { + return fmt.Errorf("checkpoint: %w", err) + } + } + return nil } diff --git a/rbf/db_test.go b/rbf/db_test.go index d637860e0..6466a3420 100644 --- a/rbf/db_test.go +++ b/rbf/db_test.go @@ -56,6 +56,70 @@ func TestDB_WAL(t *testing.T) { t.Fatalf("unexpected error: %v", err) } }) + + t.Run("Halt", func(t *testing.T) { + if testing.Short() { + t.Skip("-short enabled, skipping") + } + + config := rbfcfg.NewDefaultConfig() + config.MaxWALSize = 16 * rbf.PageSize + config.MaxWALCheckpointSize = 8 * rbf.PageSize + config.MinWALCheckpointSize = 4 * rbf.PageSize + + db := MustOpenDB(t, config) + defer MustCloseDB(t, db) + + // Continuously run read overlapping transactions. + ctx, cancel := context.WithCancel(context.Background()) + g, ctx := errgroup.WithContext(ctx) + for i := 0; i < 10; i++ { + i := i + g.Go(func() error { + time.Sleep(time.Duration(i) * 10 * time.Millisecond) // stagger + + for { + if err := ctx.Err(); err != nil { + return nil + } + + if err := func() error { + tx, err := db.Begin(false) + if err != nil { + return err + } + defer tx.Rollback() + time.Sleep(20 * time.Millisecond) + return nil + }(); err != nil { + return err + } + } + }) + } + + // Generate updates to the DB/WAL. + for i := 0; i < 100; i++ { + func() { + tx := MustBegin(t, db, true) + defer tx.Rollback() + + if err := tx.CreateBitmapIfNotExists("x"); err != nil { + t.Fatal(err) + } else if _, err := tx.Add("x", uint64(i)); err != nil { + t.Fatal(err) + } else if err := tx.Commit(); err != nil { + t.Fatal(err) + } + }() + } + + // Stop read transactions & wait. + cancel() + if err := g.Wait(); err != nil { + t.Fatal(err) + } + }) } func TestDB_Recovery(t *testing.T) { @@ -218,7 +282,6 @@ func TestDB_MultiTx(t *testing.T) { } cfg := rbfcfg.NewDefaultConfig() - cfg.CheckpointEveryDur = 1 * time.Millisecond db := MustOpenDB(t, cfg) defer MustCloseDB(t, db)