From ce656bbcdac9ab01cee45b92caebcb7b3cfcebc7 Mon Sep 17 00:00:00 2001 From: Matt Jaffee Date: Mon, 25 Feb 2019 16:49:16 -0600 Subject: [PATCH] factor out code to import/clear by positions --- fragment.go | 106 +++++++++++++++++++++++++++++---------------- roaring/roaring.go | 32 +++++++++----- 2 files changed, 89 insertions(+), 49 deletions(-) diff --git a/fragment.go b/fragment.go index 0486904e7..0eeeb8012 100644 --- a/fragment.go +++ b/fragment.go @@ -1468,23 +1468,6 @@ func (f *fragment) bulkImport(rowIDs, columnIDs []uint64, options *ImportOptions // bulkImportStandard performs a bulk import on a standard fragment. May mutate // its rowIDs and columnIDs arguments. func (f *fragment) bulkImportStandard(rowIDs, columnIDs []uint64, options *ImportOptions) (err error) { - var localBitmap *roaring.Bitmap - smallWrite := false - f.mu.Lock() - defer f.mu.Unlock() // TODO the big import path doesn't need to acquire the lock this soon. - if len(columnIDs)+f.opN < f.MaxOpN { - smallWrite = true - // If the number of operations is small, we'll write directly to local storage - localBitmap = f.storage - } else { - // Create a temporary bitmap which will be populated by rowIDs and columnIDs - // and then merged into the existing fragment's bitmap. - localBitmap = roaring.NewBitmap() - - // Disconnect op writer so we don't append updates. - localBitmap.OpWriter = nil - } - // rowSet maintains the set of rowIDs present in this import. It allows the // cache to be updated once per row, instead of once per bit. TODO: consider // sorting by rowID/columnID first and avoiding the map allocation here. (we @@ -1508,36 +1491,83 @@ func (f *fragment) bulkImportStandard(rowIDs, columnIDs []uint64, options *Impor } } positions := columnIDs - var changedN int + f.mu.Lock() + defer f.mu.Unlock() if options.Clear { - changedN, err = localBitmap.RemoveN(positions...) // TODO benchmark AddN behavior with sorted/unsorted positions + err = f.importPositions(nil, positions, rowSet) } else { - changedN, err = localBitmap.AddN(positions...) // TODO benchmark RemoveN behavior with sorted/unsorted positions + err = f.importPositions(positions, nil, rowSet) } + return errors.Wrap(err, "bulkImportStandard") +} + +// importPositions takes slices of positions within the fragment to set and +// clear in storage. One must also pass in the set of unique rows which are +// affected by the set and clear operations. It is unprotected (f.mu must be +// locked when calling it). No position should appear in both set and clear. +// +// importPositions tries to intelligently decide whether or not to do a full +// snapshot of the fragment or just do in-memory updates while appending +// operations to the op log. +func (f *fragment) importPositions(set, clear []uint64, rowSet map[uint64]struct{}) error { + var setBitmap *roaring.Bitmap + var clearBitmap *roaring.Bitmap + smallWrite := false + if len(set)+len(clear)+f.opN < f.MaxOpN { + smallWrite = true + // If the number of operations is small, we'll write directly to local storage + setBitmap = f.storage + clearBitmap = f.storage + } else { + // Create a temporary bitmap which will be populated by rowIDs and columnIDs + // and then merged into the existing fragment's bitmap. + setBitmap = roaring.NewBitmap() + clearBitmap = roaring.NewBitmap() + + // Disconnect op writer so we don't append updates. + setBitmap.OpWriter = nil + clearBitmap.OpWriter = nil + } + + changedN, err := setBitmap.AddN(set...) // TODO benchmark Add/RemoveN behavior with sorted/unsorted positions if err != nil { return errors.Wrap(err, "adding positions") } - f.stats.Count("ImportBit", int64(changedN), 1) - f.opN += changedN + + if smallWrite { + f.stats.Count("ImportBit", int64(changedN), 1) + f.opN += changedN + changedN, err = clearBitmap.RemoveN(clear...) + f.stats.Count("ClearBit", int64(changedN), 1) + f.opN += changedN + } else { + // when doing big imports, we're going to difference the local bitmap, + // so we need to add rather than removing. TODO: figure out why we + // aren't adding/removing directly from storage. + _, err = clearBitmap.AddN(clear...) + } + if err != nil { + return errors.Wrap(err, "clearing positions") + } var results *roaring.Bitmap - if !smallWrite { - // Merge localBitmap into fragment's existing data. - if options.Clear { - if f.storage.Count() > 0 { - results = f.storage.Difference(localBitmap) - } else { - results = roaring.NewBitmap() - } - } else { - if f.storage.Count() > 0 { - results = f.storage.Union(localBitmap) - } else { - results = localBitmap - } - } - } else { + if smallWrite { results = f.storage + } else if beforeCnt := f.storage.Count(); !smallWrite && beforeCnt > 0 { + // Merge localBitmap into fragment's existing data. + if len(clear) > 0 { + f.storage = f.storage.Difference(clearBitmap) + nc := f.storage.Count() + f.stats.Count("ClearBit", int64(beforeCnt-nc), 1) + beforeCnt = nc + results = f.storage + } + if len(set) > 0 { + results = f.storage.Union(setBitmap) // TODO replace with UnionInPlace after https://github.com/pilosa/pilosa/issues/1875 + f.stats.Count("ImportBit", int64(beforeCnt-f.storage.Count()), 1) + } + } else { // !smallWrite and nothing in storage + results = setBitmap } // Update cache counts for all affected rows. diff --git a/roaring/roaring.go b/roaring/roaring.go index e97803812..c8eb42016 100644 --- a/roaring/roaring.go +++ b/roaring/roaring.go @@ -131,7 +131,7 @@ func NewBitmap(a ...uint64) *Bitmap { b := &Bitmap{ Containers: newSliceContainers(), } - b.Add(a...) + b.AddN(a...) return b } @@ -177,12 +177,17 @@ func (b *Bitmap) Add(a ...uint64) (changed bool, err error) { // AddN adds values to the bitmap, appending them all to the op log in a batched write. It returns the number of changed bits. func (b *Bitmap) AddN(a ...uint64) (changed int, err error) { - op := &op{ - typ: opTypeAddBatch, - values: a, + if len(a) == 0 { + return 0, nil } - if err := b.writeOp(op); err != nil { - return 0, errors.Wrap(err, "writing to op log") + if b.OpWriter != nil { + op := &op{ + typ: opTypeAddBatch, + values: a, + } + if err := b.writeOp(op); err != nil { + return 0, errors.Wrap(err, "writing to op log") + } } // TODO consider applying changes in-memory first and then only writing the @@ -233,12 +238,17 @@ func (b *Bitmap) Remove(a ...uint64) (changed bool, err error) { } func (b *Bitmap) RemoveN(a ...uint64) (changed int, err error) { - op := &op{ - typ: opTypeRemoveBatch, - values: a, + if len(a) == 0 { + return 0, nil } - if err := b.writeOp(op); err != nil { - return 0, errors.Wrap(err, "writing to op log") + if b.OpWriter != nil { + op := &op{ + typ: opTypeRemoveBatch, + values: a, + } + if err := b.writeOp(op); err != nil { + return 0, errors.Wrap(err, "writing to op log") + } } // TODO consider applying changes in-memory first and then only writing the