Add WAL write cache mutex; update name; add benchmarks

This commit is contained in:
Ben Johnson 2020-08-18 08:41:30 -06:00
parent 5849794a1b
commit 51504f4fe4
2 changed files with 74 additions and 18 deletions

View file

@ -18,6 +18,7 @@ import (
"fmt"
"os"
"path/filepath"
"sync"
"syscall"
"github.com/pilosa/pilosa/v2/syswrap"
@ -25,12 +26,13 @@ import (
// WALSegment represents a single file in the WAL.
type WALSegment struct {
minWALID int64 // base WALID; calculated from path
path string // path to file
w *os.File // write handle
data []byte // read-only mmap data
buf []byte // write buffer
pageN int // number of written pages
mu sync.RWMutex
minWALID int64 // base WALID; calculated from path
path string // path to file
w *os.File // write handle
data []byte // read-only mmap data
writeCache []byte // write buffer
pageN int // number of written pages
}
// NewWALSegment returns a new instance of WALSegment for a given path.
@ -139,6 +141,9 @@ func (s *WALSegment) CloseForWrite() error {
// ReadWALPage reads a single page at the given WAL ID.
func (s *WALSegment) ReadWALPage(walID int64) ([]byte, error) {
s.mu.RLock()
defer s.mu.RUnlock()
// Ensure requested ID is contained in this file.
if walID < s.minWALID || walID > s.minWALID+int64(s.pageN) {
return nil, fmt.Errorf("wal segment page read out of range: id=%d base=%d pageN=%d", walID, s.minWALID, s.pageN)
@ -147,9 +152,9 @@ func (s *WALSegment) ReadWALPage(walID int64) ([]byte, error) {
offset := (walID - s.minWALID) * PageSize
// If offset is within write buffer, return from write buffer.
writeBufferOffset := int64((s.pageN * PageSize) - len(s.buf))
writeBufferOffset := int64((s.pageN * PageSize) - len(s.writeCache))
if offset >= writeBufferOffset {
buf := s.buf[offset-writeBufferOffset:]
buf := s.writeCache[offset-writeBufferOffset:]
return buf[:PageSize:PageSize], nil
}
@ -161,6 +166,9 @@ func (s *WALSegment) ReadWALPage(walID int64) ([]byte, error) {
func (s *WALSegment) WriteWALPage(page []byte, isMeta bool) (walID int64, err error) {
assert(len(page) == PageSize, "invalid page size: %d", len(page))
s.mu.Lock()
defer s.mu.Unlock()
// Initialize write file handle if not yet initialized.
if s.w == nil {
if s.w, err = os.OpenFile(s.path, os.O_WRONLY, 0666); err != nil {
@ -178,24 +186,32 @@ func (s *WALSegment) WriteWALPage(page []byte, isMeta bool) (walID int64, err er
}
// Append write to write buffer & increment page count.
s.buf = append(s.buf, page...)
s.writeCache = append(s.writeCache, page...)
s.pageN++
return walID, nil
}
// Sync flushes all changes to disk.
// Flush flushes the write buffer to the OS cache.
func (s *WALSegment) Flush() error {
s.mu.Lock()
defer s.mu.Unlock()
if _, err := s.w.WriteAt(s.writeCache, int64((s.pageN*PageSize)-len(s.writeCache))); err != nil {
return fmt.Errorf("wal segment write: %w", err)
}
s.writeCache = nil
return nil
}
// Sync flushes the write buffer and invokes a file sync to flush data to disk.
func (s *WALSegment) Sync() error {
if s.w == nil {
return nil
}
// Flush buffer to disk.
if _, err := s.w.WriteAt(s.buf, int64((s.pageN*PageSize)-len(s.buf))); err != nil {
return fmt.Errorf("wal segment write: %w", err)
if err := s.Flush(); err != nil {
return err
}
s.buf = nil
return s.w.Sync()
}

View file

@ -16,7 +16,7 @@ package rbf_test
import (
"bytes"
"encoding/hex"
"encoding/hex"
"io/ioutil"
"math/rand"
"os"
@ -26,7 +26,6 @@ import (
"github.com/pilosa/pilosa/v2/rbf"
)
func TestWALSegment_Open(t *testing.T) {
t.Run("OK", func(t *testing.T) {
s := MustOpenWALSegment(t, 10)
@ -108,6 +107,47 @@ func TestParseWALSegmentPath(t *testing.T) {
})
}
func BenchmarkWALSegment_WriteWALPage(b *testing.B) {
b.Run("8KB", func(b *testing.B) { benchmarkWALSegment_WriteWALPage(b, 8*(1<<10)) })
b.Run("16KB", func(b *testing.B) { benchmarkWALSegment_WriteWALPage(b, 16*(1<<10)) })
b.Run("64KB", func(b *testing.B) { benchmarkWALSegment_WriteWALPage(b, 64*(1<<10)) })
b.Run("256KB", func(b *testing.B) { benchmarkWALSegment_WriteWALPage(b, 256*(1<<10)) })
b.Run("1MB", func(b *testing.B) { benchmarkWALSegment_WriteWALPage(b, (1 << 20)) })
b.Run("10MB", func(b *testing.B) { benchmarkWALSegment_WriteWALPage(b, 10*(1<<20)) })
}
func benchmarkWALSegment_WriteWALPage(b *testing.B, flushSize int) {
page := make([]byte, rbf.PageSize)
for i := 0; i < b.N; i++ {
func() {
s := MustOpenWALSegment(b, 0)
defer MustCloseWALSegment(b, s)
// Fill the segment but stop after each flush interval to flush the write buffer.
for j := 0; j < rbf.MaxWALSegmentFileSize; j += rbf.PageSize {
if _, err := s.WriteWALPage(page, false); err != nil {
b.Fatal(err)
}
// Flush write buffer.
if j != 0 && j%flushSize == 0 {
if err := s.Flush(); err != nil {
b.Fatal(err)
}
}
}
// Fsync to disk at the end.
if err := s.Sync(); err != nil {
b.Fatal(err)
}
}()
}
b.SetBytes(rbf.MaxWALSegmentFileSize)
}
// MustOpenWALSegment opens a WAL segment in a temporary path. Fails on error.
func MustOpenWALSegment(tb testing.TB, walID int64) *rbf.WALSegment {
tb.Helper()