From ef156f617229aecefef8751c3fb215cb98113367 Mon Sep 17 00:00:00 2001 From: Jason Aten Date: Tue, 18 Aug 2020 16:55:44 -0500 Subject: [PATCH 1/4] ubuntu 20:10 image instead of alpine, for cgo support --- Dockerfile | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) diff --git a/Dockerfile b/Dockerfile index a30097e1e..a0950ffed 100644 --- a/Dockerfile +++ b/Dockerfile @@ -7,11 +7,13 @@ COPY . pilosa RUN cd pilosa && make install FLAGS="-a -mod=vendor ${BUILD_FLAGS}" ${MAKE_FLAGS} -FROM alpine:3.9.4 +FROM ubuntu:20.10 LABEL maintainer "dev@pilosa.com" -RUN apk add --no-cache curl jq +RUN apt-get update +## debug image: RUN apt-get install -y curl htop vim golang tree jq netcat +RUN apt-get install -y curl jq COPY --from=builder /go/bin/pilosa /pilosa From 5849794a1b4d6ec961319e63ebaedae0d2371d1a Mon Sep 17 00:00:00 2001 From: Ben Johnson Date: Fri, 14 Aug 2020 08:58:59 -0600 Subject: [PATCH 2/4] Add RBF WAL write cache --- rbf/wal.go | 29 +++++++++++++++++++++++++---- 1 file changed, 25 insertions(+), 4 deletions(-) diff --git a/rbf/wal.go b/rbf/wal.go index 0f8461e16..fae0dc713 100644 --- a/rbf/wal.go +++ b/rbf/wal.go @@ -29,6 +29,7 @@ type WALSegment struct { 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 } @@ -121,6 +122,12 @@ func (s *WALSegment) Close() error { // CloseForWrite closes the write handle, if initialized. func (s *WALSegment) CloseForWrite() error { + // Ensure write buffer is flushed out. + if err := s.Sync(); err != nil { + return err + } + + // Close underlying file writer. if s.w != nil { if err := s.w.Close(); err != nil { return err @@ -138,6 +145,15 @@ 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)) + if offset >= writeBufferOffset { + buf := s.buf[offset-writeBufferOffset:] + return buf[:PageSize:PageSize], nil + } + + // Otherwise return from on-disk mmap. return s.data[offset : offset+PageSize], nil } @@ -161,10 +177,8 @@ func (s *WALSegment) WriteWALPage(page []byte, isMeta bool) (walID int64, err er // TODO: Write meta page checksum } - // Write page at position & increment page count. - if _, err := s.w.WriteAt(page, int64(s.pageN*PageSize)); err != nil { - return 0, fmt.Errorf("wal segment write: %w", err) - } + // Append write to write buffer & increment page count. + s.buf = append(s.buf, page...) s.pageN++ return walID, nil @@ -175,6 +189,13 @@ 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) + } + s.buf = nil + return s.w.Sync() } From 51504f4fe4aa3809202aaebe1ec177d4fab84316 Mon Sep 17 00:00:00 2001 From: Ben Johnson Date: Tue, 18 Aug 2020 08:41:30 -0600 Subject: [PATCH 3/4] Add WAL write cache mutex; update name; add benchmarks --- rbf/wal.go | 48 ++++++++++++++++++++++++++++++++---------------- rbf/wal_test.go | 44 ++++++++++++++++++++++++++++++++++++++++++-- 2 files changed, 74 insertions(+), 18 deletions(-) diff --git a/rbf/wal.go b/rbf/wal.go index fae0dc713..b0fa9fbd0 100644 --- a/rbf/wal.go +++ b/rbf/wal.go @@ -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() } diff --git a/rbf/wal_test.go b/rbf/wal_test.go index 4511578c8..b41e7e792 100644 --- a/rbf/wal_test.go +++ b/rbf/wal_test.go @@ -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() From 9dbb82cf3e72ce8c9c93d1bbe3d891db82f5cfc7 Mon Sep 17 00:00:00 2001 From: Ben Johnson Date: Wed, 19 Aug 2020 08:42:53 -0600 Subject: [PATCH 4/4] WAL mutex fixes --- rbf/wal.go | 47 +++++++++++++++++++++++++++++++++++++++++------ 1 file changed, 41 insertions(+), 6 deletions(-) diff --git a/rbf/wal.go b/rbf/wal.go index b0fa9fbd0..aee4b07a3 100644 --- a/rbf/wal.go +++ b/rbf/wal.go @@ -46,20 +46,37 @@ func NewWALSegment(path string) *WALSegment { func (s *WALSegment) Path() string { return s.path } // MinWALID returns the initial WAL ID of the segment. Only available after Open(). -func (s *WALSegment) MinWALID() int64 { return s.minWALID } +func (s *WALSegment) MinWALID() int64 { + s.mu.RLock() + defer s.mu.RUnlock() + return s.minWALID +} // MaxWALID returns the maximum WAL ID of the segment. Only available after Open(). func (s *WALSegment) MaxWALID() int64 { + s.mu.RLock() + defer s.mu.RUnlock() return s.minWALID + int64(s.pageN) - 1 } // PageN returns the number of pages in the segment. -func (s *WALSegment) PageN() int { return s.pageN } +func (s *WALSegment) PageN() int { + s.mu.RLock() + defer s.mu.RUnlock() + return s.pageN +} // Size returns the current size of the segment, in bytes. -func (s *WALSegment) Size() int64 { return int64(s.pageN) * PageSize } +func (s *WALSegment) Size() int64 { + s.mu.RLock() + defer s.mu.RUnlock() + return int64(s.pageN) * PageSize +} func (s *WALSegment) Open() (err error) { + s.mu.Lock() + defer s.mu.Unlock() + // Extract base WAL ID and validate path. if s.minWALID, err = ParseWALSegmentPath(s.path); err != nil { return err @@ -110,7 +127,10 @@ func (s *WALSegment) Open() (err error) { // Close closes the write handle and the read-only mmap. func (s *WALSegment) Close() error { - if err := s.CloseForWrite(); err != nil { + s.mu.Lock() + defer s.mu.Unlock() + + if err := s.closeForWrite(); err != nil { return err } if s.data != nil { @@ -124,8 +144,14 @@ func (s *WALSegment) Close() error { // CloseForWrite closes the write handle, if initialized. func (s *WALSegment) CloseForWrite() error { + s.mu.Lock() + defer s.mu.Unlock() + return s.closeForWrite() +} + +func (s *WALSegment) closeForWrite() error { // Ensure write buffer is flushed out. - if err := s.Sync(); err != nil { + if err := s.sync(); err != nil { return err } @@ -196,7 +222,10 @@ func (s *WALSegment) WriteWALPage(page []byte, isMeta bool) (walID int64, err er func (s *WALSegment) Flush() error { s.mu.Lock() defer s.mu.Unlock() + return s.flush() +} +func (s *WALSegment) flush() error { if _, err := s.w.WriteAt(s.writeCache, int64((s.pageN*PageSize)-len(s.writeCache))); err != nil { return fmt.Errorf("wal segment write: %w", err) } @@ -206,10 +235,16 @@ func (s *WALSegment) Flush() error { // Sync flushes the write buffer and invokes a file sync to flush data to disk. func (s *WALSegment) Sync() error { + s.mu.Lock() + defer s.mu.Unlock() + return s.sync() +} + +func (s *WALSegment) sync() error { if s.w == nil { return nil } - if err := s.Flush(); err != nil { + if err := s.flush(); err != nil { return err } return s.w.Sync()