diff --git a/http/translator.go b/http/translator.go index f95067f70..0722dcbf9 100644 --- a/http/translator.go +++ b/http/translator.go @@ -22,19 +22,26 @@ import ( "io" "io/ioutil" "net/http" + "sync" "github.com/pilosa/pilosa/v2" "github.com/pilosa/pilosa/v2/logger" ) func GetOpenTranslateReaderFunc(client *http.Client) pilosa.OpenTranslateReaderFunc { + return GetOpenTranslateReaderWithLockerFunc(client, nopLocker{}) +} + +func GetOpenTranslateReaderWithLockerFunc(client *http.Client, locker sync.Locker) pilosa.OpenTranslateReaderFunc { return func(ctx context.Context, nodeURL string, offsets pilosa.TranslateOffsetMap) (pilosa.TranslateEntryReader, error) { - return openTranslateReader(ctx, nodeURL, offsets, client) + return openTranslateReader(ctx, nodeURL, offsets, client, locker) } } -func openTranslateReader(ctx context.Context, nodeURL string, offsets pilosa.TranslateOffsetMap, client *http.Client) (pilosa.TranslateEntryReader, error) { +func openTranslateReader(ctx context.Context, nodeURL string, offsets pilosa.TranslateOffsetMap, client *http.Client, locker sync.Locker) (pilosa.TranslateEntryReader, error) { r := NewTranslateEntryReader(ctx, client) + r.locker = locker + r.URL = nodeURL + "/internal/translate/data" r.Offsets = offsets if err := r.Open(); err != nil { @@ -43,9 +50,16 @@ func openTranslateReader(ctx context.Context, nodeURL string, offsets pilosa.Tra return r, nil } +type nopLocker struct{} + +func (nopLocker) Lock() {} +func (nopLocker) Unlock() {} + // TranslateEntryReader represents an implementation of pilosa.TranslateEntryReader. // It consolidates all index & field translate entries into a single reader. type TranslateEntryReader struct { + locker sync.Locker + ctx context.Context cancel func() @@ -70,7 +84,7 @@ func NewTranslateEntryReader(ctx context.Context, client *http.Client) *Translat if client == nil { client = http.DefaultClient } - r := &TranslateEntryReader{HTTPClient: client, Logger: logger.NopLogger} + r := &TranslateEntryReader{locker: nopLocker{}, HTTPClient: client, Logger: logger.NopLogger} r.ctx, r.cancel = context.WithCancel(ctx) return r } @@ -116,7 +130,10 @@ func (r *TranslateEntryReader) Close() error { r.cancel() } if r.body != nil { - return r.body.Close() + r.locker.Lock() + err := r.body.Close() + r.locker.Unlock() + return err } return nil } @@ -124,5 +141,8 @@ func (r *TranslateEntryReader) Close() error { // ReadEntry reads the next entry from the stream into entry. // Returns io.EOF at the end of the stream. func (r *TranslateEntryReader) ReadEntry(entry *pilosa.TranslateEntry) error { + r.locker.Lock() + defer r.locker.Unlock() + return r.dec.Decode(&entry) } diff --git a/http/translator_test.go b/http/translator_test.go index be6db423a..7038c1773 100644 --- a/http/translator_test.go +++ b/http/translator_test.go @@ -17,6 +17,7 @@ package http_test import ( "context" "fmt" + "sync" "testing" "time" @@ -151,3 +152,84 @@ func TestTranslateStore_EntryReader(t *testing.T) { }) */ } + +func benchmarkSetup(b *testing.B, ctx context.Context, key string, nkeys int) (string, pilosa.TranslateOffsetMap, func()) { + b.Helper() + + cluster := test.MustRunCluster(b, 1) + primary := cluster[0] + + idx := primary.MustCreateIndex(b, "i", pilosa.IndexOptions{}) + fld := primary.MustCreateField(b, idx.Name(), "f", pilosa.OptFieldKeys()) + offset := make(pilosa.TranslateOffsetMap) + offset.SetIndexPartitionOffset(idx.Name(), 0, 1) + offset.SetFieldOffset(idx.Name(), fld.Name(), 1) + + // Set data on the primary node. + for k := 0; k < nkeys; k++ { + if _, err := primary.API.Query(ctx, &pilosa.QueryRequest{ + Index: idx.Name(), + Query: fmt.Sprintf(`Set(%d, %s="%s%[1]d")`, k, fld.Name(), key), + }); err != nil { + b.Fatalf("quering api: %+v", err) + } + } + + return primary.URL(), offset, func() { + b.Helper() + + if err := primary.API.DeleteIndex(ctx, idx.Name()); err != nil { + panic(err) + } + if err := cluster.Close(); err != nil { + panic(err) + } + } +} + +func benchmarkReadEntry(b *testing.B, r pilosa.TranslateEntryReader, key string, nkeys int) { + var entry pilosa.TranslateEntry + for k := 0; k < nkeys; k++ { + if err := r.ReadEntry(&entry); err != nil { + b.Fatalf("reading entry: %+v", err) + } + if entry.Key != fmt.Sprintf("%s%d", key, k) { + b.Fatalf("got: %s, expected: %s%d", entry.Key, key, k) + } + } +} + +const ( + key = "foo" + nkeys = 1000 +) + +func BenchmarkReadEntryNoMutex(b *testing.B) { + ctx := context.Background() + url, offset, teardown := benchmarkSetup(b, ctx, key, nkeys) + defer teardown() + + for n := 0; n < b.N; n++ { + r, err := http.GetOpenTranslateReaderFunc(nil)(ctx, url, offset) + if err != nil { + b.Fatalf("openining translate reader: %+v", err) + } + benchmarkReadEntry(b, r, key, nkeys) + r.Close() + } +} + +func BenchmarkReadEntryWithMutex(b *testing.B) { + ctx := context.Background() + url, offset, teardown := benchmarkSetup(b, ctx, key, nkeys) + defer teardown() + + for n := 0; n < b.N; n++ { + r, err := http.GetOpenTranslateReaderWithLockerFunc(nil, &sync.Mutex{})(ctx, url, offset) + if err != nil { + b.Fatalf("openining translate reader: %+v", err) + } + benchmarkReadEntry(b, r, key, nkeys) + r.Close() + } +} diff --git a/server/handler_test.go b/server/handler_test.go index 31954eff8..ccd91bdcc 100644 --- a/server/handler_test.go +++ b/server/handler_test.go @@ -27,6 +27,7 @@ import ( "net/http/httptest" "reflect" "strings" + "sync" "testing" "time" @@ -1300,7 +1301,7 @@ func TestCluster_TranslateStore(t *testing.T) { cluster[0] = test.NewCommandNode(true, server.OptCommandServerOptions( pilosa.OptServerOpenTranslateStore(boltdb.OpenTranslateStore), - pilosa.OptServerOpenTranslateReader(http.GetOpenTranslateReaderFunc(nil)), + pilosa.OptServerOpenTranslateReader(http.GetOpenTranslateReaderWithLockerFunc(nil, &sync.Mutex{})), ), ) cluster[0].Config.Gossip.Port = "0" @@ -1329,7 +1330,7 @@ func TestClusterTranslator(t *testing.T) { cluster[1] = test.NewCommandNode(false, server.OptCommandServerOptions( pilosa.OptServerOpenTranslateStore(boltdb.OpenTranslateStore), - pilosa.OptServerOpenTranslateReader(http.GetOpenTranslateReaderFunc(nil)), + pilosa.OptServerOpenTranslateReader(http.GetOpenTranslateReaderWithLockerFunc(nil, &sync.Mutex{})), ), ) cluster[1].Config.Gossip.Port = "0" diff --git a/server/server.go b/server/server.go index 5578cf35f..e74d460e2 100644 --- a/server/server.go +++ b/server/server.go @@ -31,6 +31,7 @@ import ( "os/signal" "runtime" "strconv" + "sync" "syscall" "time" @@ -389,7 +390,7 @@ func (m *Command) SetupServer() error { pilosa.OptServerDiagnosticsInterval(diagnosticsInterval), pilosa.OptServerExecutorPoolSize(m.Config.WorkerPoolSize), pilosa.OptServerOpenTranslateStore(boltdb.OpenTranslateStore), - pilosa.OptServerOpenTranslateReader(http.GetOpenTranslateReaderFunc(c)), + pilosa.OptServerOpenTranslateReader(http.GetOpenTranslateReaderWithLockerFunc(c, &sync.Mutex{})), pilosa.OptServerLogger(m.logger), pilosa.OptServerAttrStoreFunc(boltdb.NewAttrStore), pilosa.OptServerSystemInfo(gopsutil.NewSystemInfo()),