Merge branch 'master' into reopen

This commit is contained in:
Kuba Podgórski 2019-05-03 01:11:15 +02:00 • committed by GitHub
commit 42c62187cf
No known key found for this signature in database
GPG key ID: 4AEE18F83AFDEB23
7 changed files with 84 additions and 41 deletions

5
api.go
View file

@ -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")
}

View file

@ -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
}

View file

@ -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.

View file

@ -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)
}

5
go.mod
View file

@ -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

5
go.sum
View file

@ -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=

View file

@ -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 {