diff --git a/api.go b/api.go index 8b30053fb..237593083 100644 --- a/api.go +++ b/api.go @@ -290,6 +290,7 @@ func setUpImportOptions(opts ...ImportOption) (*ImportOptions, error) { // bitmap. func (api *API) ImportRoaring(ctx context.Context, indexName, fieldName string, shard uint64, remote bool, req *ImportRoaringRequest) (err error) { span, ctx := tracing.StartSpanFromContext(ctx, "API.ImportRoaring") + span.LogKV("index", indexName, "field", fieldName) defer span.Finish() if err = api.validate(apiField); err != nil { @@ -325,7 +326,7 @@ func (api *API) ImportRoaring(ctx context.Context, indexName, fieldName string, } fileMagic := uint32(binary.LittleEndian.Uint16(viewData[0:2])) if fileMagic == roaring.MagicNumber { // if pilosa roaring format - err = field.importRoaring(viewData, shard, viewName, req.Clear) + err = field.importRoaring(ctx, viewData, shard, viewName, req.Clear) if err != nil { return errors.Wrap(err, "importing pilosa roaring") } @@ -335,7 +336,7 @@ func (api *API) ImportRoaring(ctx context.Context, indexName, fieldName string, // field.importRoaring changes the standard roaring run format to pilosa roaring data := make([]byte, len(viewData)) copy(data, viewData) - err = field.importRoaring(data, shard, viewName, req.Clear) + err = field.importRoaring(ctx, data, shard, viewName, req.Clear) if err != nil { return errors.Wrap(err, "importing standard roaring") } diff --git a/field.go b/field.go index 6ef362728..3b0b63feb 100644 --- a/field.go +++ b/field.go @@ -16,6 +16,7 @@ package pilosa import ( "bufio" + "context" "encoding/json" "fmt" "io/ioutil" @@ -32,6 +33,7 @@ import ( "github.com/pilosa/pilosa/pql" "github.com/pilosa/pilosa/roaring" "github.com/pilosa/pilosa/stats" + "github.com/pilosa/pilosa/tracing" "github.com/pkg/errors" ) @@ -1218,10 +1220,14 @@ func (f *Field) importValue(columnIDs []uint64, values []int64, options *ImportO return nil } -func (f *Field) importRoaring(data []byte, shard uint64, viewName string, clear bool) error { +func (f *Field) importRoaring(ctx context.Context, data []byte, shard uint64, viewName string, clear bool) error { + span, ctx := tracing.StartSpanFromContext(ctx, "Field.importRoaring") + defer span.Finish() + if viewName == "" { viewName = viewStandard } + span.LogKV("view", viewName, "bytes", len(data), "shard", shard) view, err := f.createViewIfNotExists(viewName) if err != nil { return errors.Wrap(err, "creating view") @@ -1232,7 +1238,7 @@ func (f *Field) importRoaring(data []byte, shard uint64, viewName string, clear return errors.Wrap(err, "creating fragment") } - if err := frag.importRoaring(data, clear); err != nil { + if err := frag.importRoaring(ctx, data, clear); err != nil { return err } diff --git a/fragment.go b/fragment.go index 7c24a5d66..8e8fa0ce5 100644 --- a/fragment.go +++ b/fragment.go @@ -1811,11 +1811,17 @@ func (f *fragment) importValue(columnIDs, values []uint64, bitDepth uint, clear // importRoaring imports from the official roaring data format defined at // https://github.com/RoaringBitmap/RoaringFormatSpec or from pilosa's version // of the roaring format. The cache is updated to reflect the new data. -func (f *fragment) importRoaring(data []byte, clear bool) error { +func (f *fragment) importRoaring(ctx context.Context, data []byte, clear bool) error { + span, ctx := tracing.StartSpanFromContext(ctx, "fragment.importRoaring") + defer span.Finish() + span, ctx = tracing.StartSpanFromContext(ctx, "importRoaring.AcquireFragmentLock") f.mu.Lock() defer f.mu.Unlock() + span.Finish() bm := roaring.NewBTreeBitmap() + span, ctx = tracing.StartSpanFromContext(ctx, "importRoaring.UnmarshalBinary") err := bm.UnmarshalBinary(data) + span.Finish() if err != nil { return err } @@ -1851,16 +1857,25 @@ func (f *fragment) importRoaring(data []byte, clear bool) error { if clear { toSet, toClear = toClear, toSet } - return f.importPositions(toSet, toClear, rowSet) + span, _ = tracing.StartSpanFromContext(ctx, "importRoaring.ImportPositions") + err := f.importPositions(toSet, toClear, rowSet) + span.Finish() + return err } if clear { + span, ctx = tracing.StartSpanFromContext(ctx, "importRoaringDifference") bm = f.storage.Difference(bm) + span.Finish() } else if f.storage.Containers.Size() >= bm.Containers.Size() { + span, ctx = tracing.StartSpanFromContext(ctx, "importRoaringStorageUIP") f.storage.UnionInPlace(bm) bm = f.storage + span.Finish() } else { + span, ctx = tracing.StartSpanFromContext(ctx, "importRoaringBitmapUIP") bm.UnionInPlace(f.storage) + span.Finish() } for rowID := range rowSet { @@ -1869,7 +1884,10 @@ func (f *fragment) importRoaring(data []byte, clear bool) error { } f.cache.Recalculate() - err = unprotectedWriteToFragment(f, bm) + span, _ = tracing.StartSpanFromContext(ctx, "importRoaring.WriteToFragment") + n, err := unprotectedWriteToFragment(f, bm) + span.LogKV("bytesWritten", n) + span.Finish() return err } @@ -1900,12 +1918,13 @@ func track(start time.Time, message string, stats stats.StatsClient, logger logg } func (f *fragment) snapshot() error { - return unprotectedWriteToFragment(f, f.storage) + _, err := unprotectedWriteToFragment(f, f.storage) + return err } // unprotectedWriteToFragment writes the fragment f with bm as the data. It is unprotected, and // f.mu must be locked when calling it. -func unprotectedWriteToFragment(f *fragment, bm *roaring.Bitmap) error { // nolint: interfacer +func unprotectedWriteToFragment(f *fragment, bm *roaring.Bitmap) (n int64, err error) { // nolint: interfacer completeMessage := fmt.Sprintf("fragment: snapshot complete %s/%s/%s/%d", f.index, f.field, f.view, f.shard) start := time.Now() @@ -1915,39 +1934,39 @@ func unprotectedWriteToFragment(f *fragment, bm *roaring.Bitmap) error { // noli snapshotPath := f.path + snapshotExt file, err := os.Create(snapshotPath) if err != nil { - return fmt.Errorf("create snapshot file: %s", err) + return n, fmt.Errorf("create snapshot file: %s", err) } defer file.Close() // Write storage to snapshot. bw := bufio.NewWriter(file) - if _, err := bm.WriteTo(bw); err != nil { - return fmt.Errorf("snapshot write to: %s", err) + if n, err = bm.WriteTo(bw); err != nil { + return n, fmt.Errorf("snapshot write to: %s", err) } if err := bw.Flush(); err != nil { - return fmt.Errorf("flush: %s", err) + return n, fmt.Errorf("flush: %s", err) } // Close current storage. if err := f.closeStorage(); err != nil { - return fmt.Errorf("close storage: %s", err) + return n, fmt.Errorf("close storage: %s", err) } // Move snapshot to data file location. if err := os.Rename(snapshotPath, f.path); err != nil { - return fmt.Errorf("rename snapshot: %s", err) + return n, fmt.Errorf("rename snapshot: %s", err) } // Reopen storage. if err := f.openStorage(); err != nil { - return fmt.Errorf("open storage: %s", err) + return n, fmt.Errorf("open storage: %s", err) } // Reset operation count. f.opN = 0 - return nil + return n, nil } // RecalculateCache rebuilds the cache regardless of invalidate time delay. diff --git a/fragment_internal_test.go b/fragment_internal_test.go index db0159734..43b11357c 100644 --- a/fragment_internal_test.go +++ b/fragment_internal_test.go @@ -16,6 +16,7 @@ package pilosa import ( "bytes" + "context" "flag" "fmt" "io" @@ -745,7 +746,7 @@ func BenchmarkFragment_RepeatedSmallImports(b *testing.B) { f := mustOpenFragment("i", "f", viewStandard, 0, "") f.MaxOpN = opN defer f.Clean(b) - err := f.importRoaring(getZipfRowsSliceRoaring(uint64(numRows), 1, 0, ShardWidth), false) + err := f.importRoaringT(getZipfRowsSliceRoaring(uint64(numRows), 1, 0, ShardWidth), false) if err != nil { b.Fatalf("importing base data for benchmark: %v", err) } @@ -781,14 +782,14 @@ func BenchmarkFragment_RepeatedSmallImportsRoaring(b *testing.B) { f := mustOpenFragment("i", "f", viewStandard, 0, "") f.MaxOpN = opN defer f.Clean(b) - err := f.importRoaring(getZipfRowsSliceRoaring(numRows, 1, 0, ShardWidth), false) + err := f.importRoaringT(getZipfRowsSliceRoaring(numRows, 1, 0, ShardWidth), false) if err != nil { b.Fatalf("importing base data for benchmark: %v", err) } for i := 0; i < numUpdates; i++ { data := getUpdataRoaring(numRows, bitsPerUpdate, int64(i)) b.StartTimer() - err := f.importRoaring(data, false) + err := f.importRoaringT(data, false) b.StopTimer() if err != nil { b.Fatalf("doing small roaring import: %v", err) @@ -1024,7 +1025,7 @@ func TestFragment_TopN_Intersect_Large(t *testing.T) { if err != nil { t.Fatalf("writing to bytes: %v", err) } - err = f.importRoaring(b.Bytes(), false) + err = f.importRoaringT(b.Bytes(), false) if err != nil { t.Fatalf("importing data: %v", err) } @@ -2008,7 +2009,7 @@ func BenchmarkImportRoaring(b *testing.B) { for i := 0; i < b.N; i++ { f := mustOpenFragment("i", fmt.Sprintf("r%dc%s", numRows, cacheType), viewStandard, 0, cacheType) b.StartTimer() - err := f.importRoaring(data, false) + err := f.importRoaringT(data, false) if err != nil { f.Clean(b) b.Fatalf("import error: %v", err) @@ -2041,7 +2042,7 @@ func BenchmarkImportRoaringConcurrent(b *testing.B) { for j := 0; j < concurrency; j++ { j := j eg.Go(func() error { - return frags[j].importRoaring(data, false) + return frags[j].importRoaringT(data, false) }) } err := eg.Wait() @@ -2072,7 +2073,7 @@ func BenchmarkImportRoaringUpdateConcurrent(b *testing.B) { for i := 0; i < b.N; i++ { for j := 0; j < concurrency; j++ { frags[j] = mustOpenFragment("i", "f", viewStandard, uint64(j), CacheTypeRanked) - err := frags[j].importRoaring(data, false) + err := frags[j].importRoaringT(data, false) if err != nil { b.Fatalf("importing roaring: %v", err) } @@ -2082,7 +2083,7 @@ func BenchmarkImportRoaringUpdateConcurrent(b *testing.B) { for j := 0; j < concurrency; j++ { j := j eg.Go(func() error { - return frags[j].importRoaring(updata, false) + return frags[j].importRoaringT(updata, false) }) } err := eg.Wait() @@ -2138,12 +2139,12 @@ func BenchmarkImportRoaringUpdate(b *testing.B) { b.StopTimer() for i := 0; i < b.N; i++ { f := mustOpenFragment("i", fmt.Sprintf("r%dc%s", numRows, cacheType), viewStandard, 0, cacheType) - err := f.importRoaring(data, false) + err := f.importRoaringT(data, false) if err != nil { b.Errorf("import error: %v", err) } b.StartTimer() - err = f.importRoaring(updata, false) + err = f.importRoaringT(updata, false) if err != nil { f.Clean(b) b.Errorf("import error: %v", err) @@ -2174,12 +2175,12 @@ func BenchmarkUpdatePathological(b *testing.B) { for i := 0; i < b.N; i++ { b.StopTimer() f := mustOpenFragment("i", "f", viewStandard, 0, DefaultCacheType) - err := f.importRoaring(exists, false) + err := f.importRoaringT(exists, false) if err != nil { b.Fatalf("importing roaring: %v", err) } b.StartTimer() - err = f.importRoaring(inc, false) + err = f.importRoaringT(inc, false) if err != nil { b.Fatalf("importing second: %v", err) } @@ -2196,7 +2197,7 @@ func initBigFrag() { for i := int64(0); i < 10; i++ { // 10 million rows, 1 bit per column, random seeded by i data := getZipfRowsSliceRoaring(10000000, i, 0, ShardWidth) - err := f.importRoaring(data, false) + err := f.importRoaringT(data, false) if err != nil { panic(fmt.Sprintf("setting up fragment data: %v", err)) } @@ -2273,7 +2274,7 @@ func BenchmarkImportRoaringIntoLargeFragment(b *testing.B) { b.Fatalf("opening fragment: %v", err) } b.StartTimer() - err = nf.importRoaring(updata, false) + err = nf.importRoaringT(updata, false) b.StopTimer() if err != nil { b.Fatalf("bulkImport: %v", err) @@ -2286,7 +2287,7 @@ func BenchmarkImportRoaringIntoLargeFragment(b *testing.B) { func TestGetZipfRowsSliceRoaring(t *testing.T) { f := mustOpenFragment("i", "f", viewStandard, 0, DefaultCacheType) data := getZipfRowsSliceRoaring(10, 1, 0, ShardWidth) - err := f.importRoaring(data, false) + err := f.importRoaringT(data, false) if err != nil { t.Fatalf("importing roaring: %v", err) } @@ -2431,6 +2432,11 @@ func (f *fragment) Clean(t testing.TB) { } } +// importRoaringT calls importRoaring with context.Background() for convenience +func (f *fragment) importRoaringT(data []byte, clear bool) error { + return f.importRoaring(context.Background(), data, clear) +} + // CleanKeep is just like Clean(), but it doesn't remove the // fragment file (note that it DOES remove the cache file). func (f *fragment) CleanKeep(t testing.TB) { @@ -2622,7 +2628,7 @@ func TestFragment_RoaringImport(t *testing.T) { if err != nil { t.Fatalf("writing to buffer: %v", err) } - err = f.importRoaring(buf.Bytes(), false) + err = f.importRoaringT(buf.Bytes(), false) if err != nil { t.Fatalf("importing roaring: %v", err) } @@ -2697,7 +2703,7 @@ func TestFragment_RoaringImportTopN(t *testing.T) { if err != nil { t.Fatalf("writing to buffer: %v", err) } - err = f.importRoaring(buf.Bytes(), false) + err = f.importRoaringT(buf.Bytes(), false) if err != nil { t.Fatalf("importing roaring: %v", err) } @@ -2941,7 +2947,7 @@ func TestUnionInPlaceMapped(t *testing.T) { count1 := setBM1.Count() // now we write setBM0 into f.storage. - err = unprotectedWriteToFragment(f, setBM0) + _, err = unprotectedWriteToFragment(f, setBM0) if err != nil { t.Fatalf("trying to flush fragment to disk: %v", err) } diff --git a/go.mod b/go.mod index 4234b36fa..397d0e267 100644 --- a/go.mod +++ b/go.mod @@ -29,9 +29,8 @@ require ( github.com/spf13/cobra v0.0.3 github.com/spf13/pflag v1.0.3 github.com/spf13/viper v1.3.1 - github.com/uber-go/atomic v1.3.2 // indirect - github.com/uber/jaeger-client-go v2.15.0+incompatible - github.com/uber/jaeger-lib v1.5.0 + github.com/uber/jaeger-client-go v2.16.0+incompatible + github.com/uber/jaeger-lib v2.0.0+incompatible // indirect golang.org/x/crypto v0.0.0-20190426145343-a29dc8fdc734 // indirect golang.org/x/net v0.0.0-20190424112056-4829fb13d2c6 // indirect golang.org/x/sync v0.0.0-20190423024810-112230192c58 diff --git a/go.sum b/go.sum index 7fb7bfd6f..a1541510b 100644 --- a/go.sum +++ b/go.sum @@ -54,6 +54,8 @@ github.com/hashicorp/golang-lru v0.5.0 h1:CL2msUPvZTLb5O648aiLNJw3hnBxN2+1Jq8rCO github.com/hashicorp/golang-lru v0.5.0/go.mod h1:/m3WP610KZHVQ1SGc6re/UDhFvYD7pJ4Ao+sR/qLZy8= github.com/hashicorp/hcl v1.0.0 h1:0Anlzjpi4vEasTeNFn2mLJgTSwt0+6sfsiTG8qcWGx4= github.com/hashicorp/hcl v1.0.0/go.mod h1:E5yfLk+7swimpb2L/Alb/PJmXilQ/rhwaUYs4T20WEQ= +github.com/hashicorp/memberlist v0.1.3 h1:EmmoJme1matNzb+hMpDuR/0sbJSUisxyqBGG676r31M= +github.com/hashicorp/memberlist v0.1.3/go.mod h1:ajVTdAv/9Im8oMAAj5G31PhhMCZJV2pPBoIllUwCN7I= github.com/inconshreveable/mousetrap v1.0.0 h1:Z8tu5sraLXCXIcARxBp/8cbvlwVa7Z1NHg9XEKhtSvM= github.com/inconshreveable/mousetrap v1.0.0/go.mod h1:PxqpIevigyE2G7u3NXJIT2ANytuPF1OarO4DADm73n8= github.com/magiconair/properties v1.8.0 h1:LLgXmsheXeRoUOBOjtwPQCWIYqM/LU1ayDtDePerRcY= @@ -84,7 +86,6 @@ github.com/shirou/gopsutil v2.18.12+incompatible h1:1eaJvGomDnH74/5cF4CTmTbLHAri github.com/shirou/gopsutil v2.18.12+incompatible/go.mod h1:5b4v6he4MtMOwMlS0TUMTu2PcXUg8+E1lC7eC3UO/RA= github.com/shirou/w32 v0.0.0-20160930032740-bb4de0191aa4 h1:udFKJ0aHUL60LboW/A+DfgoHVedieIzIXE8uylPue0U= github.com/shirou/w32 v0.0.0-20160930032740-bb4de0191aa4/go.mod h1:qsXQc7+bwAM3Q1u/4XEfrquwF8Lw7D7y5cD8CuHnfIc= -github.com/spaolacci/murmur3 v0.0.0-20180118202830-f09979ecbc72 h1:qLC7fQah7D6K1B0ujays3HV9gkFtllcxhzImRR7ArPQ= github.com/spaolacci/murmur3 v0.0.0-20180118202830-f09979ecbc72/go.mod h1:JwIasOWyU6f++ZhiEuf87xNszmSA2myDM2Kzu9HwQUA= github.com/spf13/afero v1.1.2 h1:m8/z1t7/fwjysjQRYbP0RD+bUIF/8tJwPdEZsI83ACI= github.com/spf13/afero v1.1.2/go.mod h1:j4pytiNVoe2o6bmDsKpLACNPDBIoEAkihy7loJ1B0CQ= @@ -104,6 +105,8 @@ github.com/uber-go/atomic v1.3.2 h1:Azu9lPBWRNKzYXSIwRfgRuDuS0YKsK4NFhiQv98gkxo= github.com/uber-go/atomic v1.3.2/go.mod h1:/Ct5t2lcmbJ4OSe/waGBoaVvVqtO0bmtfVNex1PFV8g= github.com/uber/jaeger-client-go v2.15.0+incompatible h1:NP3qsSqNxh8VYr956ur1N/1C1PjvOJnJykCzcD5QHbk= github.com/uber/jaeger-client-go v2.15.0+incompatible/go.mod h1:WVhlPFC8FDjOFMMWRy2pZqQJSXxYSwNYOkTr/Z6d3Kk= +github.com/uber/jaeger-client-go v2.16.0+incompatible h1:Q2Pp6v3QYiocMxomCaJuwQGFt7E53bPYqEgug/AoBtY= +github.com/uber/jaeger-client-go v2.16.0+incompatible/go.mod h1:WVhlPFC8FDjOFMMWRy2pZqQJSXxYSwNYOkTr/Z6d3Kk= github.com/uber/jaeger-lib v1.5.0 h1:OHbgr8l656Ub3Fw5k9SWnBfIEwvoHQ+W2y+Aa9D1Uyo= github.com/uber/jaeger-lib v1.5.0/go.mod h1:ComeNDZlWwrWnDv8aPp0Ba6+uUTzImX/AauajbLI56U= github.com/ugorji/go/codec v0.0.0-20181204163529-d75b2dcb6bc8/go.mod h1:VFNgLljTbGfSG7qAOspJ7OScBnGdDN/yBr0sguwnwf0= diff --git a/http/handler.go b/http/handler.go index b52ee52a3..78b5f48a9 100644 --- a/http/handler.go +++ b/http/handler.go @@ -1617,15 +1617,23 @@ func (h *Handler) handlePostImportRoaring(w http.ResponseWriter, r *http.Request remote = true } + ctx := r.Context() + // Read entire body. + span, _ := tracing.StartSpanFromContext(ctx, "ioutil.ReadAll-Body") body, err := ioutil.ReadAll(r.Body) + span.LogKV("bodySize", len(body)) + span.Finish() if err != nil { http.Error(w, err.Error(), http.StatusBadRequest) return } req := &pilosa.ImportRoaringRequest{} - if err := h.api.Serializer.Unmarshal(body, req); err != nil { + span, _ = tracing.StartSpanFromContext(ctx, "Unmarshal") + err = h.api.Serializer.Unmarshal(body, req) + span.Finish() + if err != nil { http.Error(w, err.Error(), http.StatusBadRequest) return } @@ -1639,7 +1647,7 @@ func (h *Handler) handlePostImportRoaring(w http.ResponseWriter, r *http.Request resp := &pilosa.ImportResponse{} // TODO give meaningful stats for import - err = h.api.ImportRoaring(r.Context(), indexName, fieldName, shard, remote, req) + err = h.api.ImportRoaring(ctx, indexName, fieldName, shard, remote, req) if err != nil { resp.Err = err.Error() if _, ok := err.(pilosa.BadRequestError); ok { @@ -1648,6 +1656,7 @@ func (h *Handler) handlePostImportRoaring(w http.ResponseWriter, r *http.Request w.WriteHeader(http.StatusInternalServerError) } } + // Marshal response object. buf, err := h.api.Serializer.Marshal(resp) if err != nil {