thread the holder through things, and improve snapshot queue logic

This is logically two separate things, but the individual changes
are thoroughly intertwined in the code.

The first change is a logical change to the design of the snapshot
queue, which is that it now adjusts the maxOpN the background scan
targets, allowing it to lower that value over time when things are
quiet. We do this because it turns out that on large data sets,
this can make a factor-of-four difference in memory usage!

So, in general, on a quiet system, each pass through the holder
aims for about 1/4 of the existing fragments to get snapshotted.
When there's more load, we adjust those values up.

We also make the snapshot queue a bit less chatty, to make testing
less annoying -- we only print stats if the queue enqueues at least
two snapshots, or skips any.

The second change is threading the holder through things. We've
always threaded the logger through, and then added the snapshot
queue, and some of the Inspect-related work led to wanting to
have a way to thread options through, so what if we just threaded
the holder itself through, and removed the direct copying around
of the logger, snapshot queue, and so on. Similarly, everything
can now use holder.PartitionN instead of having to get its own
copy of PartitionN handed out to each index.

This does imply ensuring that test cases always get a reasonable
default holder.

This is a precursor to adding additional information to the holder,
such as whether it's in a special read-only mode, which would imply
not modifying on-disk files. This is already semi-supported for
the specific case of the background snapshot queue and cache flushing,
which are attached to the (created in a previous commit) new
holder Activate method, instead of being automatic on holder Open.

The change to a snapshot queue can also cause races in tests, because
the fragment.Clean method's "sanity check" accesses a fragment without
a lock. Fix that. Since there's a couple of t.Fatalf(), but we need
to release the lock before closing, we use an anonymous function
with a defer to handle that. Whee!
This commit is contained in:
Seebs 2020-04-07 18:51:04 -05:00 • committed by Jaden Weiss
parent 121717594b
commit ceaf5c15d1
No known key found for this signature in database
GPG key ID: 177F065773634B67
19 changed files with 235 additions and 195 deletions

View file

@ -96,7 +96,7 @@ func newIndexWithTempPath(name string) *Index {
if err != nil {
panic(err)
}
index, err := NewIndex(path, name, DefaultPartitionN)
index, err := NewIndex(NewHolder(DefaultPartitionN), path, name)
if err != nil {
panic(err)
}

View file

@ -31,7 +31,6 @@ import (
"github.com/gogo/protobuf/proto"
"github.com/pilosa/pilosa/v2/internal"
"github.com/pilosa/pilosa/v2/logger"
"github.com/pilosa/pilosa/v2/pql"
"github.com/pilosa/pilosa/v2/roaring"
"github.com/pilosa/pilosa/v2/stats"
@ -117,10 +116,6 @@ type Field struct {
// Shards with data on any node in the cluster, according to this node.
remoteAvailableShards *roaring.Bitmap
logger logger.Logger
snapshotQueue snapshotQueue
translateStore TranslateStore
// Instantiates new translation stores
@ -338,17 +333,17 @@ func OptFieldTypeBool() FieldOption {
// that it's of the type `OptFieldType*`). This means
// this function couldn't be used to set, for example,
// `FieldOptions.Keys`.
func NewField(path, index, name string, opts FieldOption) (*Field, error) {
func NewField(holder *Holder, path, index, name string, opts FieldOption) (*Field, error) {
err := validateName(name)
if err != nil {
return nil, errors.Wrap(err, "validating name")
}
return newField(path, index, name, opts)
return newField(holder, path, index, name, opts)
}
// newField returns a new instance of field (without name validation).
func newField(path, index, name string, opts FieldOption) (*Field, error) {
func newField(holder *Holder, path, index, name string, opts FieldOption) (*Field, error) {
// Apply functional option.
fo := FieldOptions{}
err := opts(&fo)
@ -372,7 +367,7 @@ func newField(path, index, name string, opts FieldOption) (*Field, error) {
remoteAvailableShards: roaring.NewBitmap(),
logger: logger.NopLogger,
holder: holder,
OpenTranslateStore: OpenInMemTranslateStore,
}
@ -448,7 +443,7 @@ func (f *Field) loadAvailableShards() error {
}
// some other problem:
if err != nil {
f.logger.Printf("available shards file present but unreadable, discarding: %v", err)
f.holder.Logger.Printf("available shards file present but unreadable, discarding: %v", err)
err = os.Remove(path)
if err != nil {
return errors.Wrap(err, "deleting corrupt available shards list")
@ -457,7 +452,7 @@ func (f *Field) loadAvailableShards() error {
}
bm := roaring.NewBitmap()
if err = bm.UnmarshalBinary(buf); err != nil {
f.logger.Printf("available shards file corrupt, discarding: %v", err)
f.holder.Logger.Printf("available shards file corrupt, discarding: %v", err)
err = os.Remove(path)
if err != nil {
return errors.Wrap(err, "deleting corrupt available shards list")
@ -548,17 +543,17 @@ func (f *Field) Options() FieldOptions {
func (f *Field) Open() error {
if err := func() (err error) {
// Ensure the field's path exists.
f.logger.Debugf("ensure field path exists: %s", f.path)
f.holder.Logger.Debugf("ensure field path exists: %s", f.path)
if err := os.MkdirAll(f.path, 0777); err != nil {
return errors.Wrap(err, "creating field dir")
}
f.logger.Debugf("load meta file for index/field: %s/%s", f.index, f.name)
f.holder.Logger.Debugf("load meta file for index/field: %s/%s", f.index, f.name)
if err := f.loadMeta(); err != nil {
return errors.Wrap(err, "loading meta")
}
f.logger.Debugf("load available shards for index/field: %s/%s", f.index, f.name)
f.holder.Logger.Debugf("load available shards for index/field: %s/%s", f.index, f.name)
if err := f.loadAvailableShards(); err != nil {
return errors.Wrap(err, "loading available shards")
}
@ -570,17 +565,17 @@ func (f *Field) Open() error {
}
// Apply the field options loaded from meta (or set via setOptions()).
f.logger.Debugf("apply options for index/field: %s/%s", f.index, f.name)
f.holder.Logger.Debugf("apply options for index/field: %s/%s", f.index, f.name)
if err := f.applyOptions(f.options); err != nil {
return errors.Wrap(err, "applying options")
}
f.logger.Debugf("open views for index/field: %s/%s", f.index, f.name)
f.holder.Logger.Debugf("open views for index/field: %s/%s", f.index, f.name)
if err := f.openViews(); err != nil {
return errors.Wrap(err, "opening views")
}
f.logger.Debugf("open row attribute store for index/field: %s/%s", f.index, f.name)
f.holder.Logger.Debugf("open row attribute store for index/field: %s/%s", f.index, f.name)
if err := f.rowAttrStore.Open(); err != nil {
return errors.Wrap(err, "opening attrstore")
}
@ -607,7 +602,7 @@ func (f *Field) Open() error {
return err
}
f.logger.Debugf("successfully opened field index/field: %s/%s", f.index, f.name)
f.holder.Logger.Debugf("successfully opened field index/field: %s/%s", f.index, f.name)
return nil
}
func blockingWriteAvailableShards(fieldPath string, availableShardBytes []byte) {
@ -748,7 +743,7 @@ fileLoop:
<-fieldQueue
}()
name := filepath.Base(fi.Name())
f.logger.Debugf("open index/field/view: %s/%s/%s", f.index, f.name, fi.Name())
f.holder.Logger.Debugf("open index/field/view: %s/%s/%s", f.index, f.name, fi.Name())
view := f.newView(f.viewPath(name), name)
if err := view.open(); err != nil {
return fmt.Errorf("opening view: view=%s, err=%s", view.name, err)
@ -770,7 +765,7 @@ fileLoop:
}
view.rowAttrStore = f.rowAttrStore
f.logger.Debugf("add index/field/view to field.viewMap: %s/%s/%s", f.index, f.name, view.name)
f.holder.Logger.Debugf("add index/field/view to field.viewMap: %s/%s/%s", f.index, f.name, view.name)
mu.Lock()
f.viewMap[view.name] = view
mu.Unlock()
@ -1183,14 +1178,10 @@ func (f *Field) createViewIfNotExistsBase(name string) (*view, bool, error) {
}
func (f *Field) newView(path, name string) *view {
view := newView(path, f.index, f.name, name, f.options)
view.logger = f.logger
view := newView(f.holder, path, f.index, f.name, name, f.options)
view.rowAttrStore = f.rowAttrStore
view.stats = f.Stats
view.broadcaster = f.broadcaster
if f.snapshotQueue != nil {
view.snapshotQueue = f.snapshotQueue
}
return view
}

View file

@ -205,7 +205,7 @@ func NewTestField(t *testing.T, opts FieldOption) *TestField {
if err != nil {
t.Fatal(err)
}
field, err := NewField(path, "i", "f", opts)
field, err := NewField(NewHolder(DefaultPartitionN), path, "i", "f", opts)
if err != nil {
t.Fatal(err)
}
@ -235,7 +235,7 @@ func (f *TestField) Reopen() error {
}
path, index, name := f.Path(), f.Index(), f.Name()
f.Field, err = NewField(path, index, name, OptFieldTypeDefault())
f.Field, err = NewField(NewHolder(DefaultPartitionN), path, index, name, OptFieldTypeDefault())
if err != nil {
return err
}
@ -730,7 +730,7 @@ func TestDecimalField_MinMaxBoundaries(t *testing.T) {
},
} {
t.Run("minmax"+strconv.Itoa(i), func(t *testing.T) {
_, err := NewField("no-path", "i", "f", OptFieldTypeDecimal(test.scale, test.min, test.max))
_, err := NewField(NewHolder(DefaultPartitionN), "no-path", "i", "f", OptFieldTypeDecimal(test.scale, test.min, test.max))
if err != nil && test.expErr {
if !strings.Contains(err.Error(), "is not supported") {
t.Fatal(err)

View file

@ -144,7 +144,7 @@ func TestField_NameRestriction(t *testing.T) {
if err != nil {
panic(err)
}
field, err := pilosa.NewField(path, "i", ".meta", pilosa.OptFieldTypeDefault())
field, err := pilosa.NewField(pilosa.NewHolder(pilosa.DefaultPartitionN), path, "i", ".meta", pilosa.OptFieldTypeDefault())
if field != nil {
t.Fatalf("unexpected field name %s", err)
}
@ -177,13 +177,13 @@ func TestField_NameValidation(t *testing.T) {
panic(err)
}
for _, name := range validFieldNames {
_, err := pilosa.NewField(path, "i", name, pilosa.OptFieldTypeDefault())
_, err := pilosa.NewField(pilosa.NewHolder(pilosa.DefaultPartitionN), path, "i", name, pilosa.OptFieldTypeDefault())
if err != nil {
t.Fatalf("unexpected field name: %s %s", name, err)
}
}
for _, name := range invalidFieldNames {
_, err := pilosa.NewField(path, "i", name, pilosa.OptFieldTypeDefault())
_, err := pilosa.NewField(pilosa.NewHolder(pilosa.DefaultPartitionN), path, "i", name, pilosa.OptFieldTypeDefault())
if err == nil {
t.Fatalf("expected error on field name: %s", name)
}

View file

@ -106,6 +106,10 @@ type fragment struct {
field string
view string
shard uint64
// parent holder, used to find snapshot queue, etc.
holder *Holder
// debugging tool: addresses of current and previous maps
prevdata, currdata struct{ from, to uintptr }
@ -152,12 +156,10 @@ type fragment struct {
mutexVector vector
stats stats.StatsClient
snapshotQueue snapshotQueue
}
// newFragment returns a new instance of Fragment.
func newFragment(path, index, field, view string, shard uint64, flags byte) *fragment {
func newFragment(holder *Holder, path, index, field, view string, shard uint64, flags byte) *fragment {
f := &fragment{
path: path,
index: index,
@ -168,11 +170,10 @@ func newFragment(path, index, field, view string, shard uint64, flags byte) *fra
CacheType: DefaultCacheType,
CacheSize: DefaultCacheSize,
Logger: logger.NopLogger,
holder: holder,
MaxOpN: defaultFragmentMaxOpN,
stats: stats.NopStatsClient,
snapshotQueue: defaultSnapshotQueue,
stats: stats.NopStatsClient,
}
f.snapshotCond = sync.Cond{L: &f.mu}
return f
@ -188,13 +189,13 @@ func (f *fragment) Open() error {
if err := func() error {
// Initialize storage in a function so we can close if anything goes wrong.
f.Logger.Debugf("open storage for index/field/view/fragment: %s/%s/%s/%d", f.index, f.field, f.view, f.shard)
f.holder.Logger.Debugf("open storage for index/field/view/fragment: %s/%s/%s/%d", f.index, f.field, f.view, f.shard)
if err := f.openStorage(true); err != nil {
return errors.Wrap(err, "opening storage")
}
// Fill cache with rows persisted to disk.
f.Logger.Debugf("open cache for index/field/view/fragment: %s/%s/%s/%d", f.index, f.field, f.view, f.shard)
f.holder.Logger.Debugf("open cache for index/field/view/fragment: %s/%s/%s/%d", f.index, f.field, f.view, f.shard)
if err := f.openCache(); err != nil {
e2 := f.closeStorage()
if e2 != nil {
@ -214,7 +215,7 @@ func (f *fragment) Open() error {
return err
}
f.Logger.Debugf("successfully opened index/field/view/fragment: %s/%s/%s/%d", f.index, f.field, f.view, f.shard)
f.holder.Logger.Debugf("successfully opened index/field/view/fragment: %s/%s/%s/%d", f.index, f.field, f.view, f.shard)
return nil
}
@ -273,7 +274,7 @@ func (f *fragment) importStorage(data []byte, file *os.File, newGen generation,
}
return false, fmt.Errorf("unmarshal storage: file=%s, err=%s", file.Name(), err)
}
f.Logger.Printf("warning: unmarshal storage, file=%s, err=%v", file.Name(), err)
f.holder.Logger.Printf("warning: unmarshal storage, file=%s, err=%v", file.Name(), err)
trunc, ok := cause.(roaring.FileShouldBeTruncatedError)
if ok {
// generation code looks for a FileShouldBeTruncatedError
@ -299,7 +300,7 @@ func (f *fragment) applyStorage(data []byte, file *os.File, newGen generation, m
if file != nil {
fi, err := file.Stat()
if err != nil {
f.Logger.Printf("trying to apply new storage to existing bitmap, stat failed: %v", err)
f.holder.Logger.Printf("trying to apply new storage to existing bitmap, stat failed: %v", err)
}
if err == nil && fi != nil && fi.Size() == 0 {
return f.emptyStorage(file)
@ -359,7 +360,7 @@ func (f *fragment) openStorage(unmarshalData bool) error {
storageOp = f.applyStorage
}
var err error
f.gen, err = newGeneration(f.gen, f.path, unmarshalData, storageOp, f.Logger)
f.gen, err = newGeneration(f.gen, f.path, unmarshalData, storageOp, f.holder.Logger)
if f.gen != nil {
scratchData := f.gen.Bytes()
f.prevdata = f.currdata
@ -407,7 +408,7 @@ func (f *fragment) openCache() error {
// Unmarshal cache data.
var pb internal.Cache
if err := proto.Unmarshal(buf, &pb); err != nil {
f.Logger.Printf("error unmarshaling cache data, skipping: path=%s, err=%s", path, err)
f.holder.Logger.Printf("error unmarshaling cache data, skipping: path=%s, err=%s", path, err)
return nil
}
@ -435,13 +436,13 @@ func (f *fragment) Close() error {
func (f *fragment) close() error {
// Flush cache if closing gracefully.
if err := f.flushCache(); err != nil {
f.Logger.Printf("fragment: error flushing cache on close: err=%s, path=%s", err, f.path)
f.holder.Logger.Printf("fragment: error flushing cache on close: err=%s, path=%s", err, f.path)
return errors.Wrap(err, "flushing cache")
}
// Close underlying storage.
if err := f.closeStorage(); err != nil {
f.Logger.Printf("fragment: error closing storage: err=%s, path=%s", err, f.path)
f.holder.Logger.Printf("fragment: error closing storage: err=%s, path=%s", err, f.path)
return errors.Wrap(err, "closing storage")
}
@ -688,7 +689,8 @@ func (f *fragment) unprotectedSetRow(row *Row, rowID uint64) (changed bool, err
f.rowCache.Add(rowID, nil)
// Snapshot storage.
f.snapshotQueue.Enqueue(f)
f.holder.SnapshotQueue.Enqueue(f)
f.stats.Count("setRow", 1, 1.0)
return changed, nil
}
@ -728,7 +730,7 @@ func (f *fragment) unprotectedClearRow(rowID uint64) (changed bool, err error) {
f.rowCache.Add(rowID, nil)
// Snapshot storage.
f.snapshotQueue.Enqueue(f)
f.holder.SnapshotQueue.Enqueue(f)
return changed, nil
}
@ -2006,11 +2008,11 @@ func (f *fragment) importPositions(set, clear []uint64, rowSet map[uint64]struct
// we got an error. it's possible that the error indicates that something went wrong.
mappedIn, mappedOut, unmappedIn, errs, e2 := f.storage.SanityCheckMapping(f.currdata.from, f.currdata.to)
if errs != 0 {
f.Logger.Printf("transaction failed on %s. storage has %d mapped in range, %d mapped out of range, %d unmapped in range, %d errors total, last %v",
f.holder.Logger.Printf("transaction failed on %s. storage has %d mapped in range, %d mapped out of range, %d unmapped in range, %d errors total, last %v",
f.path, mappedIn, mappedOut, unmappedIn, errs, e2)
if f.prevdata.from != f.currdata.from {
mappedIn, mappedOut, unmappedIn, errs, e2 = f.storage.SanityCheckMapping(f.prevdata.from, f.prevdata.to)
f.Logger.Printf("with previous map, storage would have %d mapped in range, %d mapped out of range, %d unmapped in range, %d errors total, last %v",
f.holder.Logger.Printf("with previous map, storage would have %d mapped in range, %d mapped out of range, %d unmapped in range, %d errors total, last %v",
mappedIn, mappedOut, unmappedIn, errs, e2)
}
}
@ -2164,7 +2166,7 @@ func (f *fragment) importValue(columnIDs []uint64, values []int64, bitDepth uint
// in theory, this should probably have been queued anyway, but if enough
// of the bits matched existing bits, we'll be under our opN estimate, and
// we want to ensure that the snapshot happens.
return f.snapshotQueue.Immediate(f)
return f.holder.SnapshotQueue.Immediate(f)
}
// importRoaring imports from the official roaring data format defined at
@ -2252,7 +2254,7 @@ func (f *fragment) incrementOpN(changed int) {
f.opN += changed
f.ops++
if f.opN > f.MaxOpN {
f.snapshotQueue.Enqueue(f)
f.holder.SnapshotQueue.Enqueue(f)
}
}
@ -2286,7 +2288,7 @@ func (f *fragment) snapshot() (err error) {
// we can't see the actual values that were used to generate this, probably.
if e2.Error() == "runtime error: invalid memory address or nil pointer dereference" {
mappedIn, mappedOut, unmappedIn, errs, _ := f.storage.SanityCheckMapping(f.currdata.from, f.currdata.to)
f.Logger.Printf("transaction failed on %s. storage has %d mapped in range, %d mapped out of range, %d unmapped in range, %d errors total",
f.holder.Logger.Printf("transaction failed on %s. storage has %d mapped in range, %d mapped out of range, %d unmapped in range, %d errors total",
f.path, mappedIn, mappedOut, unmappedIn, errs)
}
} else {
@ -2306,7 +2308,7 @@ func (f *fragment) snapshot() (err error) {
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()
defer track(start, completeMessage, f.stats, f.Logger)
defer track(start, completeMessage, f.stats, f.holder.Logger)
// Create a temporary file to snapshot to.
snapshotPath := f.path + snapshotExt
@ -3096,7 +3098,7 @@ func (s *fragmentSyncer) syncBlockFromPrimary(id int) error {
// the primary node.
nodes := s.Cluster.shardNodes(f.index, f.shard)
if s.Node.ID != nodes[0].ID {
f.Logger.Debugf("non-primary replica expecting sync from primary: %s, index=%s, field=%s, shard=%d", nodes[0].ID, f.index, f.field, f.shard)
f.holder.Logger.Debugf("non-primary replica expecting sync from primary: %s, index=%s, field=%s, shard=%d", nodes[0].ID, f.index, f.field, f.shard)
return nil
}

View file

@ -1477,7 +1477,7 @@ func BenchmarkFragment_Blocks(b *testing.B) {
}
// Open the fragment specified by the path.
f := newFragment(*FragmentPath, "i", "f", viewStandard, 0, 0)
f := newFragment(NewHolder(DefaultPartitionN), *FragmentPath, "i", "f", viewStandard, 0, 0)
if err := f.Open(); err != nil {
b.Fatal(err)
}
@ -2006,7 +2006,7 @@ func BenchmarkFragment_Snapshot(b *testing.B) {
b.ReportAllocs()
// Open the fragment specified by the path.
f := newFragment(*FragmentPath, "i", "f", viewStandard, 0, 0)
f := newFragment(NewHolder(DefaultPartitionN), *FragmentPath, "i", "f", viewStandard, 0, 0)
if err := f.Open(); err != nil {
b.Fatal(err)
}
@ -2121,7 +2121,7 @@ func BenchmarkImportRoaring(b *testing.B) {
// care whether this succeeds,
// but if it's happening we want
// it to be done.
_ = f.snapshotQueue.Await(f)
_ = defaultSnapshotQueue.Await(f)
f.Clean(b)
b.Fatalf("import error: %v", err)
}
@ -2162,7 +2162,7 @@ func BenchmarkImportRoaringConcurrent(b *testing.B) {
err := frags[j].importRoaringT(data[j], false)
// error unimportant if it happened, but we want
// any snapshots to have finished.
_ = frags[j].snapshotQueue.Await(frags[j])
_ = defaultSnapshotQueue.Await(frags[j])
return err
})
}
@ -2203,7 +2203,7 @@ func BenchmarkImportRoaringUpdateConcurrent(b *testing.B) {
if err != nil {
b.Fatalf("importing roaring: %v", err)
}
err = frags[j].snapshotQueue.Immediate(frags[j])
err = defaultSnapshotQueue.Immediate(frags[j])
if err != nil {
b.Fatalf("snapshot after import: %v", err)
}
@ -2214,7 +2214,7 @@ func BenchmarkImportRoaringUpdateConcurrent(b *testing.B) {
j := j
eg.Go(func() error {
err := frags[j].importRoaringT(updata, false)
err2 := frags[j].snapshotQueue.Await(frags[j])
err2 := defaultSnapshotQueue.Await(frags[j])
if err == nil {
err = err2
}
@ -2282,7 +2282,7 @@ func BenchmarkImportRoaringUpdate(b *testing.B) {
if err != nil {
b.Errorf("import error: %v", err)
}
err = f.snapshotQueue.Immediate(f)
err = defaultSnapshotQueue.Immediate(f)
if err != nil {
b.Errorf("snapshot after import error: %v", err)
}
@ -2292,7 +2292,7 @@ func BenchmarkImportRoaringUpdate(b *testing.B) {
f.Clean(b)
b.Errorf("import error: %v", err)
}
err = f.snapshotQueue.Await(f)
err = defaultSnapshotQueue.Await(f)
if err != nil {
b.Errorf("snapshot after import error: %v", err)
}
@ -2391,7 +2391,7 @@ func BenchmarkImportIntoLargeFragment(b *testing.B) {
}
origF.Close()
fi.Close()
nf := newFragment(fi.Name(), "i", "f", viewStandard, 0, 0)
nf := newFragment(NewHolder(DefaultPartitionN), fi.Name(), "i", "f", viewStandard, 0, 0)
err = nf.Open()
if err != nil {
b.Fatalf("opening fragment: %v", err)
@ -2428,7 +2428,7 @@ func BenchmarkImportRoaringIntoLargeFragment(b *testing.B) {
}
origF.Close()
fi.Close()
nf := newFragment(fi.Name(), "i", "f", viewStandard, 0, 0)
nf := newFragment(NewHolder(DefaultPartitionN), fi.Name(), "i", "f", viewStandard, 0, 0)
err = nf.Open()
if err != nil {
b.Fatalf("opening fragment: %v", err)
@ -2607,17 +2607,23 @@ func (f *fragment) sanityCheck(t testing.TB) {
func (f *fragment) Clean(t testing.TB) {
f.mu.Lock()
err := f.snapshotQueue.Await(f)
f.mu.Unlock()
if err != nil {
t.Fatalf("snapshot failed before sanity check: %v", err)
}
f.sanityCheck(t)
if f.storage != nil && f.storage.Source != nil {
if f.storage.Source.Dead() {
t.Fatalf("cleaning up fragment %s, source %s, source already dead", f.path, f.storage.Source.ID())
// we need to ensure that we unlock the mutex before terminating
// the clean operation, but we need it held during the sanity
// check or else, in some cases, the background snapshot queue
// can decide to pick it up.
func() {
defer f.mu.Unlock()
err := defaultSnapshotQueue.Await(f)
if err != nil {
t.Fatalf("snapshot failed before sanity check: %v", err)
}
}
f.sanityCheck(t)
if f.storage != nil && f.storage.Source != nil {
if f.storage.Source.Dead() {
t.Fatalf("cleaning up fragment %s, source %s, source already dead", f.path, f.storage.Source.ID())
}
}
}()
errc := f.Close()
// prevent double-closes of generation during testing.
f.gen = nil
@ -2626,10 +2632,6 @@ func (f *fragment) Clean(t testing.TB) {
if errc != nil || errf != nil {
t.Fatal("cleaning up fragment: ", errc, errf, errp)
}
if f.snapshotQueue != nil {
f.snapshotQueue.Stop()
f.snapshotQueue = nil
}
// not all fragments have cache files
if errp != nil && !os.IsNotExist(errp) {
t.Fatalf("cleaning up fragment cache: %v", errp)
@ -2649,10 +2651,6 @@ func (f *fragment) CleanKeep(t testing.TB) {
if errc != nil {
t.Fatal("closing fragment: ", errc, errp)
}
if f.snapshotQueue != nil {
f.snapshotQueue.Stop()
f.snapshotQueue = nil
}
// not all fragments have cache files
if errp != nil && !os.IsNotExist(errp) {
t.Fatalf("cleaning up fragment cache: %v", errp)
@ -2680,12 +2678,11 @@ func mustOpenFragmentFlags(index, field, view string, shard uint64, cacheType st
cacheType = DefaultCacheType
}
f := newFragment(file.Name(), index, field, view, shard, flags)
f := newFragment(NewHolder(DefaultPartitionN), file.Name(), index, field, view, shard, flags)
f.CacheType = cacheType
f.RowAttrStore = &memAttrStore{
store: make(map[uint64]map[string]interface{}),
}
f.snapshotQueue = newSnapshotQueue(1, 1, nil)
if err := f.Open(); err != nil {
panic(err)
@ -3179,7 +3176,7 @@ func TestUnionInPlaceMapped(t *testing.T) {
// it's used only in computation of things that usually don't go to
// disk, which is why we handle this specially in testing and not
// generically.
err = f.snapshotQueue.Immediate(f)
err = defaultSnapshotQueue.Immediate(f)
if err != nil {
t.Fatalf("snapshot after union-in-place: %v", err)
}
@ -3332,7 +3329,6 @@ func TestImportClearRestart(t *testing.T) {
cols: []uint64{1, 1, 1, 1, 1, 1},
},
}
for i, test := range tests {
for _, maxOpN := range []int{0, 10000} {
t.Run(fmt.Sprintf("%dMaxOpN%d", i, maxOpN), func(t *testing.T) {
@ -3385,7 +3381,7 @@ func TestImportClearRestart(t *testing.T) {
check(t, f, exp)
f2 := newFragment(f.path, "i", "f", viewStandard, 0, 0)
f2 := newFragment(NewHolder(DefaultPartitionN), f.path, "i", "f", viewStandard, 0, 0)
f2.MaxOpN = maxOpN
f2.CacheType = f.CacheType
@ -3419,7 +3415,7 @@ func TestImportClearRestart(t *testing.T) {
check(t, f2, exp)
f3 := newFragment(f2.path, "i", "f", viewStandard, 0, 0)
f3 := newFragment(NewHolder(DefaultPartitionN), f2.path, "i", "f", viewStandard, 0, 0)
f3.MaxOpN = maxOpN
f3.CacheType = f.CacheType

View file

@ -39,14 +39,14 @@ func TestGenerationPanic(t *testing.T) {
}
prevData = f.gen.(*mmapGeneration).data
f.mu.Lock()
_ = f.snapshotQueue.Immediate(f)
_ = defaultSnapshotQueue.Immediate(f)
f.mu.Unlock()
runtime.GC()
for i := 0; i < (f.MaxOpN / 2); i++ {
_, _ = f.setBit(0, uint64(i*32)+23)
}
f.mu.Lock()
f.snapshotQueue.Await(f)
defaultSnapshotQueue.Await(f)
f.mu.Unlock()
runtime.GC()
newData := f.gen.(*mmapGeneration).data

View file

@ -77,9 +77,8 @@ type Holder struct {
// The interval at which the cached row ids are persisted to disk.
cacheFlushInterval time.Duration
Logger logger.Logger
snapshotQueue snapshotQueue
Logger logger.Logger
SnapshotQueue SnapshotQueue
// Instantiates new translation stores
OpenTranslateStore OpenTranslateStoreFunc
@ -169,6 +168,8 @@ func NewHolder(partitionN int) *Holder {
translationSyncer: NopTranslationSyncer,
Logger: logger.NopLogger,
SnapshotQueue: defaultSnapshotQueue,
}
}
@ -213,11 +214,6 @@ func (h *Holder) Open() error {
return errors.Wrap(err, "reading directory")
}
// Run snapshots asynchronously. The snapshotQueue will have a background
// task associated with it which flushes it and waits until this channel
// is closed, so we should always close this channel when done.
h.snapshotQueue = newSnapshotQueue(10, 2, h.Logger)
for _, fi := range fis {
// Skip files or hidden directories.
if !fi.IsDir() || strings.HasPrefix(fi.Name(), ".") {
@ -261,18 +257,25 @@ func (h *Holder) Open() error {
h.Logger.Printf("open holder: complete")
// Periodically flush cache.
h.wg.Add(1)
go func() { defer h.wg.Done(); h.monitorCacheFlush() }()
h.Stats.Open()
h.snapshotQueue.ScanHolder(h)
h.opened.Close()
return nil
}
// Activate runs the background tasks relevant to keeping a holder in a stable
// state, such as scanning it for needed snapshots, or flushing caches. This
// is separate from opening because, while a server would nearly always want
// to do this, other use cases (like consistency checks of a data directory)
// need to avoid it even getting started.
func (h *Holder) Activate() {
// Periodically flush cache.
h.wg.Add(2)
go func() { defer h.wg.Done(); h.monitorCacheFlush() }()
go func() { defer h.wg.Done(); h.SnapshotQueue.ScanHolder(h, h.closing) }()
}
// checkForeignIndex is a check before applying a foreign
// index to a field; if the index is not yet available,
// (because holder is still opening and may not have opened
@ -313,10 +316,6 @@ func (h *Holder) Close() error {
return errors.Wrap(err, "closing index")
}
}
if h.snapshotQueue != nil {
h.snapshotQueue.Stop()
h.snapshotQueue = nil
}
// Reset opened in case Holder needs to be reopened.
h.opened.mu.Lock()
@ -587,19 +586,16 @@ func (h *Holder) createIndex(name string, opt IndexOptions) (*Index, error) {
}
func (h *Holder) newIndex(path, name string) (*Index, error) {
index, err := NewIndex(path, name, h.partitionN)
index, err := NewIndex(h, path, name)
if err != nil {
return nil, err
}
index.logger = h.Logger
index.Stats = h.Stats.WithTags(fmt.Sprintf("index:%s", index.Name()))
index.broadcaster = h.broadcaster
index.newAttrStore = h.NewAttrStore
index.columnAttrs = h.NewAttrStore(filepath.Join(index.path, ".data"))
index.snapshotQueue = h.snapshotQueue
index.OpenTranslateStore = h.OpenTranslateStore
index.translationSyncer = h.translationSyncer
index.holder = h
return index, nil
}

View file

@ -27,7 +27,6 @@ import (
"github.com/gogo/protobuf/proto"
"github.com/pilosa/pilosa/v2/internal"
"github.com/pilosa/pilosa/v2/logger"
"github.com/pilosa/pilosa/v2/roaring"
"github.com/pilosa/pilosa/v2/stats"
"github.com/pkg/errors"
@ -46,9 +45,6 @@ type Index struct {
trackExistence bool
existenceFld *Field
// Partitions used by translation.
partitionN int
// Fields by name.
fields map[string]*Field
@ -60,9 +56,6 @@ type Index struct {
broadcaster broadcaster
Stats stats.StatsClient
logger logger.Logger
snapshotQueue snapshotQueue
// Passed to field for foreign-index lookup.
holder *Holder
@ -76,24 +69,23 @@ type Index struct {
}
// NewIndex returns a new instance of Index.
func NewIndex(path, name string, partitionN int) (*Index, error) {
func NewIndex(holder *Holder, path, name string) (*Index, error) {
err := validateName(name)
if err != nil {
return nil, errors.Wrap(err, "validating name")
}
return &Index{
path: path,
name: name,
partitionN: partitionN,
fields: make(map[string]*Field),
path: path,
name: name,
fields: make(map[string]*Field),
newAttrStore: newNopAttrStore,
columnAttrs: nopStore,
broadcaster: NopBroadcaster,
Stats: stats.NopStatsClient,
logger: logger.NopLogger,
holder: holder,
trackExistence: true,
translateStores: make(map[int]TranslateStore),
@ -155,18 +147,18 @@ func (i *Index) OpenWithTimestamp() error { return i.open(true) }
func (i *Index) open(withTimestamp bool) (err error) {
// Ensure the path exists.
i.logger.Debugf("ensure index path exists: %s", i.path)
i.holder.Logger.Debugf("ensure index path exists: %s", i.path)
if err := os.MkdirAll(i.path, 0777); err != nil {
return errors.Wrap(err, "creating directory")
}
// Read meta file.
i.logger.Debugf("load meta file for index: %s", i.name)
i.holder.Logger.Debugf("load meta file for index: %s", i.name)
if err := i.loadMeta(); err != nil {
return errors.Wrap(err, "loading meta file")
}
i.logger.Debugf("open fields for index: %s", i.name)
i.holder.Logger.Debugf("open fields for index: %s", i.name)
if err := i.openFields(withTimestamp); err != nil {
return errors.Wrap(err, "opening fields")
}
@ -181,15 +173,15 @@ func (i *Index) open(withTimestamp bool) (err error) {
return errors.Wrap(err, "opening attrstore")
}
i.logger.Debugf("open translate store for index: %s", i.name)
i.holder.Logger.Debugf("open translate store for index: %s", i.name)
var g errgroup.Group
var mu sync.Mutex
for partitionID := 0; partitionID < i.partitionN; partitionID++ {
for partitionID := 0; partitionID < i.holder.partitionN; partitionID++ {
partitionID := partitionID
g.Go(func() error {
store, err := i.OpenTranslateStore(i.TranslateStorePath(partitionID), i.name, "", partitionID, i.partitionN)
store, err := i.OpenTranslateStore(i.TranslateStorePath(partitionID), i.name, "", partitionID, i.holder.partitionN)
if err != nil {
return errors.Wrapf(err, "opening index translate store: partition=%d", partitionID)
}
@ -239,7 +231,7 @@ fileLoop:
defer func() {
<-indexQueue
}()
i.logger.Debugf("open field: %s", fi.Name())
i.holder.Logger.Debugf("open field: %s", fi.Name())
mu.Lock()
fld, err := i.newField(i.fieldPath(filepath.Base(fi.Name())), filepath.Base(fi.Name()))
if withTimestamp {
@ -257,7 +249,7 @@ fileLoop:
if err := fld.Open(); err != nil {
return fmt.Errorf("open field: name=%s, err=%s", fld.Name(), err)
}
i.logger.Debugf("add field to index.fields: %s", fi.Name())
i.holder.Logger.Debugf("add field to index.fields: %s", fi.Name())
mu.Lock()
i.fields[fld.Name()] = fld
mu.Unlock()
@ -512,17 +504,13 @@ func (i *Index) createField(name string, opt *FieldOptions) (*Field, error) {
}
func (i *Index) newField(path, name string) (*Field, error) {
f, err := newField(path, i.name, name, OptFieldTypeDefault())
f, err := newField(i.holder, path, i.name, name, OptFieldTypeDefault())
if err != nil {
return nil, err
}
f.logger = i.logger
f.Stats = i.Stats
f.broadcaster = i.broadcaster
f.rowAttrStore = i.newAttrStore(filepath.Join(f.path, ".data"))
if i.snapshotQueue != nil {
f.snapshotQueue = i.snapshotQueue
}
f.OpenTranslateStore = i.OpenTranslateStore
return f, nil
}

View file

@ -25,7 +25,7 @@ func mustOpenIndex(opt IndexOptions) *Index {
if err != nil {
panic(err)
}
index, err := NewIndex(path, "i", DefaultPartitionN)
index, err := NewIndex(NewHolder(1), path, "i")
if err != nil {
panic(err)
}

View file

@ -242,7 +242,7 @@ func TestIndex_InvalidName(t *testing.T) {
if err != nil {
panic(err)
}
index, err := pilosa.NewIndex(path, "ABC", pilosa.DefaultPartitionN)
index, err := pilosa.NewIndex(pilosa.NewHolder(pilosa.DefaultPartitionN), path, "ABC")
if err == nil {
t.Fatalf("should have gotten an error on index name with caps")
}

View file

@ -2513,8 +2513,8 @@ type BitmapInfo struct {
OpDetails []OpInfo
BitCount uint64
ContainerCount int
Containers []ContainerInfo
OpContainers []ContainerInfo
Containers []ContainerInfo // The containers found in the bitmap originally
OpContainers []ContainerInfo // The containers resulting from ops log changes.
}
// Iterator represents an iterator over a Bitmap.

View file

@ -96,8 +96,9 @@ func (b *Bitmap) UnmarshalBinary(data []byte) (err error) {
return nil
}
// unmarshalPilosaRoaring treats data as being encoded in Pilosa's 64 bit
// roaring format and decodes it into b.
// InspectBinary reads a roaring bitmap, plus a possible ops log,
// and reports back on the contents, including distinguishing between
// the original ops log and the post-ops-log contents.
func InspectBinary(data []byte) (info BitmapInfo, err error) {
if data == nil {
return info, errors.New("no roaring bitmap provided")

View file

@ -66,9 +66,10 @@ type Server struct { // nolint: maligned
extensions []*ext.ExtensionInfo
// External
systemInfo SystemInfo
gcNotifier GCNotifier
logger logger.Logger
systemInfo SystemInfo
gcNotifier GCNotifier
logger logger.Logger
snapshotQueue SnapshotQueue
nodeID string
uri URI
@ -533,6 +534,9 @@ func (s *Server) UpAndDown() error {
func (s *Server) Open() error {
s.logger.Printf("open server")
// Start background monitoring.
s.snapshotQueue = newSnapshotQueue(10, 2, s.logger)
// Log startup
err := s.holder.logStartup()
if err != nil {
@ -560,6 +564,8 @@ func (s *Server) Open() error {
if err := s.holder.Open(); err != nil {
return errors.Wrap(err, "opening Holder")
}
// bring up the background tasks for the holder.
s.holder.Activate()
if err := s.cluster.setNodeState(nodeStateReady); err != nil {
return errors.Wrap(err, "setting nodeState")
}
@ -571,7 +577,6 @@ func (s *Server) Open() error {
// buffered channel.
s.cluster.listenForJoins()
// Start background monitoring.
s.wg.Add(3)
go func() { defer s.wg.Done(); s.monitorAntiEntropy() }()
go func() { defer s.wg.Done(); s.monitorRuntime() }()
@ -596,6 +601,11 @@ func (s *Server) Close() error {
if s.holder != nil {
errh = s.holder.Close()
}
if s.snapshotQueue != nil {
s.holder.SnapshotQueue = nil
s.snapshotQueue.Stop()
s.snapshotQueue = nil
}
// prefer to return holder error over cluster
// error. This order is somewhat arbitrary. It would be better if we had
// some way to combine all the errors, but probably not important enough to

View file

@ -16,6 +16,7 @@ package pilosa
import (
"fmt"
"math/bits"
"os"
"sync"
"sync/atomic"
@ -39,12 +40,22 @@ import (
// Await, Enqueue, and Immediate should be called only with the fragment lock
// held.
//
// ScanHolder spawns a new goroutine. You don't need to use `go` on it.
type snapshotQueue interface {
// If you create a queue, it should get stopped at some point. The
// atomicSnapshotQueue implementation used as defaultSnapshotQueue has
// a Start function which will tell you whether it actually started a
// queue. This logic exists because in a normal server case, you probably
// want the queue to be shut down as part of server shutdown, but if you're
// running cluster tests, you probably want to start and shop the queue as
// part of the test, not stop it when any server terminates.
//
// It's less likely to be desireable to start/stop individual queues,
// because fragments use the defaultSnapshotQueue anyway. This design
// needs revisiting.
type SnapshotQueue interface {
Immediate(*fragment) error
Enqueue(*fragment)
Await(*fragment) error
ScanHolder(*Holder)
ScanHolder(*Holder, chan struct{})
Stop()
}
@ -64,20 +75,18 @@ func (q *queuelessSnapshotQueue) Immediate(f *fragment) error {
return f.snapshot()
}
func (q *queuelessSnapshotQueue) ScanHolder(h *Holder) {
func (q *queuelessSnapshotQueue) ScanHolder(h *Holder, done chan struct{}) {
}
func (q *queuelessSnapshotQueue) Stop() {
}
// defaultSnapshotQueue is the fallback to use if none is available,
// and currently uses queueless -- it runs all snapshots immediately.
var defaultSnapshotQueue *queuelessSnapshotQueue
var defaultSnapshotQueue = &queuelessSnapshotQueue{}
// newSnapshotQueue makes a new snapshot queue, of depth N, with
// w worker threads.
func newSnapshotQueue(n int, w int, l logger.Logger) snapshotQueue {
sq := prioritySnapshotQueue{normal: make(chan snapshotRequest, n), urgent: make(chan snapshotRequest), background: make(chan snapshotRequest), done: make(chan struct{}), logger: l}
func newSnapshotQueue(n int, w int, l logger.Logger) SnapshotQueue {
sq := prioritySnapshotQueue{normal: make(chan snapshotRequest, n), urgent: make(chan snapshotRequest), background: make(chan snapshotRequest), done: make(chan struct{}), maxOpN: 10000, logger: l}
if sq.logger == nil {
sq.logger = logger.NewStandardLogger(os.Stderr)
}
@ -108,9 +117,11 @@ type prioritySnapshotQueue struct {
done chan struct{}
mu sync.RWMutex
scanWG, workerWG sync.WaitGroup
maxOpN int
observedOpN [16]int
stats struct {
enqueued uint64
skipped uint64
enqueued uint32
skipped uint32
}
}
@ -192,7 +203,9 @@ func (sq *prioritySnapshotQueue) Stop() {
sq.urgent = nil
close(sq.background)
sq.background = nil
if sq.stats.skipped > 0 || sq.stats.enqueued > 1 {
enqueued := atomic.LoadUint32(&sq.stats.enqueued)
skipped := atomic.LoadUint32(&sq.stats.skipped)
if skipped > 0 || enqueued > 1 {
sq.logger.Printf("snapshot queue: enqueued %d, skipped %d\n", sq.stats.enqueued, sq.stats.skipped)
}
}
@ -217,10 +230,10 @@ func (sq *prioritySnapshotQueue) Enqueue(f *fragment) {
// try to enqueue snapshot
select {
case sq.normal <- snapshotRequest{frag: f, when: time.Now()}:
atomic.AddUint64(&sq.stats.enqueued, 1)
atomic.AddUint32(&sq.stats.enqueued, 1)
return
default:
atomic.AddUint64(&sq.stats.skipped, 1)
atomic.AddUint32(&sq.stats.skipped, 1)
f.snapshotPending = false
return
}
@ -230,6 +243,11 @@ func (sq *prioritySnapshotQueue) Enqueue(f *fragment) {
// held. Await waits on a condition variable inside f, associated with the
// fragment's lock, so this does not conflict with the lock being used for
// snapshots.
//
// Note that workers don't stop just because the queue's been stopped; only
// the background scanner is stopped. So an Await shouldn't block forever
// even if the queue gets shut down. If you're reading this, possibly that
// analysis is incorrect.
func (sq *prioritySnapshotQueue) Await(f *fragment) (err error) {
for f.snapshotPending {
f.snapshotCond.Wait()
@ -281,9 +299,19 @@ func (sq *prioritySnapshotQueue) needsSnapshot(f *fragment) bool {
if f.snapshotPending {
return false
}
if f.opN > f.MaxOpN {
if f.opN > sq.maxOpN {
return true
}
// aka "log2(n) + 1", or 0 for n==0
pow2 := 32 - bits.LeadingZeros32(uint32(f.opN))
// 15 == 16384. we assume that since 16384 is higher than our
// normal maxOpN, it's always a reasonable value.
if pow2 > 15 {
pow2 = 15
}
// store in inverse order so the lowest slot in the array is the
// highest cardinality
sq.observedOpN[15-pow2]++
return false
}
@ -291,10 +319,10 @@ func (sq *prioritySnapshotQueue) needsSnapshot(f *fragment) bool {
// indexes/fields/views/fragments, looking for fragments which have OpN
// high enough to justify a snapshot but don't seem to have one pending.
// It then dumps these in the low priority background queue.
func (sq *prioritySnapshotQueue) ScanHolder(h *Holder) {
func (sq *prioritySnapshotQueue) ScanHolder(h *Holder, done chan struct{}) {
sq.mu.Lock()
sq.scanWG.Add(1)
go sq.scanHolderWorker(h, sq.background, sq.done)
go sq.scanHolderWorker(h, sq.background, done)
sq.mu.Unlock()
}
@ -302,6 +330,9 @@ func (sq *prioritySnapshotQueue) ScanHolder(h *Holder) {
// fragments which need snapshots taken. It's the cleanup task for snapshots
// that would have been requested by Enqueue, but the queue was full.
func (sq *prioritySnapshotQueue) scanHolderWorker(h *Holder, background chan snapshotRequest, done chan struct{}) {
// queueDone is global to this snapshotQueue, done is specific to this holder
// scanner.
queueDone := sq.done
defer sq.scanWG.Done()
var indexNames, fieldNames, viewNames []string
var fragNums []uint64
@ -379,9 +410,11 @@ func (sq *prioritySnapshotQueue) scanHolderWorker(h *Holder, background chan sna
// background queue. If there's nothing that needs snapshots,
// we pause frequently for a second or so at a time.
counter++
if counter == 100 {
if counter == 1000 {
select {
case <-time.After(1 * time.Second):
case <-queueDone:
return
case <-done:
return
}
@ -400,9 +433,36 @@ func (sq *prioritySnapshotQueue) scanHolderWorker(h *Holder, background chan sna
// No reason to be active if we're not finding anything.
select {
case <-time.After(60 * time.Second):
case <-queueDone:
return
case <-done:
return
}
}
sq.logger.Debugf("observedOpN by power of 2: %d\n", sq.observedOpN[:])
total := 0
for _, v := range sq.observedOpN {
total += v
}
target := total / 4
subTotal := 0
for i, v := range sq.observedOpN {
subTotal += v
if subTotal >= target {
prevMaxOpN := sq.maxOpN
sq.maxOpN = (1 << (15 - uint(i))) / 2
if sq.maxOpN > 0 {
sq.maxOpN--
}
if prevMaxOpN != sq.maxOpN {
sq.logger.Printf("background scan: %d/%d fragments considered have opN %d or higher\n",
subTotal, total, sq.maxOpN)
}
break
}
}
for i := range sq.observedOpN {
sq.observedOpN[i] = 0
}
}
}

View file

@ -33,7 +33,7 @@ func newField(opts pilosa.FieldOption) *Field {
if err != nil {
panic(err)
}
field, err := pilosa.NewField(path, "i", "f", opts)
field, err := pilosa.NewField(pilosa.NewHolder(pilosa.DefaultPartitionN), path, "i", "f", opts)
if err != nil {
panic(err)
}
@ -63,7 +63,7 @@ func (f *Field) reopen() error {
}
path, index, name := f.Path(), f.Index(), f.Name()
f.Field, err = pilosa.NewField(path, index, name, pilosa.OptFieldTypeDefault())
f.Field, err = pilosa.NewField(pilosa.NewHolder(pilosa.DefaultPartitionN), path, index, name, pilosa.OptFieldTypeDefault())
if err != nil {
return err
}

View file

@ -32,7 +32,7 @@ func newIndex() *Index {
if err != nil {
panic(err)
}
index, err := pilosa.NewIndex(path, "i", pilosa.DefaultPartitionN)
index, err := pilosa.NewIndex(pilosa.NewHolder(pilosa.DefaultPartitionN), path, "i")
if err != nil {
panic(err)
}
@ -62,7 +62,7 @@ func (i *Index) Reopen() error {
}
path, name := i.Path(), i.Name()
i.Index, err = pilosa.NewIndex(path, name, pilosa.DefaultPartitionN)
i.Index, err = pilosa.NewIndex(pilosa.NewHolder(pilosa.DefaultPartitionN), path, name)
if err != nil {
return err
}

42
view.go
View file

@ -26,7 +26,6 @@ import (
"sync/atomic"
"time"
"github.com/pilosa/pilosa/v2/logger"
"github.com/pilosa/pilosa/v2/pql"
"github.com/pilosa/pilosa/v2/roaring"
"github.com/pilosa/pilosa/v2/stats"
@ -49,6 +48,8 @@ type view struct {
field string
name string
holder *Holder
fieldType string
cacheType string
cacheSize uint32
@ -56,24 +57,24 @@ type view struct {
// Fragments by shard.
fragments map[uint64]*fragment
broadcaster broadcaster
stats stats.StatsClient
rowAttrStore AttrStore
logger logger.Logger
snapshotQueue snapshotQueue
broadcaster broadcaster
stats stats.StatsClient
rowAttrStore AttrStore
knownShards *roaring.Bitmap
knownShardsCopied uint32
}
// newView returns a new instance of View.
func newView(path, index, field, name string, fieldOptions FieldOptions) *view {
func newView(holder *Holder, path, index, field, name string, fieldOptions FieldOptions) *view {
return &view{
path: path,
index: index,
field: field,
name: name,
holder: holder,
fieldType: fieldOptions.Type,
cacheType: fieldOptions.CacheType,
cacheSize: fieldOptions.CacheSize,
@ -82,7 +83,6 @@ func newView(path, index, field, name string, fieldOptions FieldOptions) *view {
broadcaster: NopBroadcaster,
stats: stats.NopStatsClient,
logger: logger.NopLogger,
knownShards: roaring.NewSliceBitmap(),
}
}
@ -133,14 +133,14 @@ func (v *view) open() error {
if err := func() error {
// Ensure the view's path exists.
v.logger.Debugf("ensure view path exists: %s", v.path)
v.holder.Logger.Debugf("ensure view path exists: %s", v.path)
if err := os.MkdirAll(v.path, 0777); err != nil {
return errors.Wrap(err, "creating view directory")
} else if err := os.MkdirAll(filepath.Join(v.path, "fragments"), 0777); err != nil {
return errors.Wrap(err, "creating fragments directory")
}
v.logger.Debugf("open fragments for index/field/view: %s/%s/%s", v.index, v.field, v.name)
v.holder.Logger.Debugf("open fragments for index/field/view: %s/%s/%s", v.index, v.field, v.name)
if err := v.openFragments(); err != nil {
return errors.Wrap(err, "opening fragments")
}
@ -151,7 +151,7 @@ func (v *view) open() error {
return err
}
v.logger.Debugf("successfully opened index/field/view: %s/%s/%s", v.index, v.field, v.name)
v.holder.Logger.Debugf("successfully opened index/field/view: %s/%s/%s", v.index, v.field, v.name)
return nil
}
@ -190,12 +190,12 @@ fileLoop:
// Parse filename into integer.
shard, err := strconv.ParseUint(filepath.Base(fi.Name()), 10, 64)
if err != nil {
v.logger.Debugf("WARNING: couldn't use non-integer file as shard in index/field/view %s/%s/%s: %s", v.index, v.field, v.name, fi.Name())
v.holder.Logger.Debugf("WARNING: couldn't use non-integer file as shard in index/field/view %s/%s/%s: %s", v.index, v.field, v.name, fi.Name())
continue
}
workQueue <- struct{}{}
v.logger.Debugf("open index/field/view/fragment: %s/%s/%s/%d", v.index, v.field, v.name, shard)
v.holder.Logger.Debugf("open index/field/view/fragment: %s/%s/%s/%d", v.index, v.field, v.name, shard)
eg.Go(func() error {
defer func() {
<-workQueue
@ -205,7 +205,7 @@ fileLoop:
return fmt.Errorf("open fragment: shard=%d, err=%s", frag.shard, err)
}
frag.RowAttrStore = v.rowAttrStore
v.logger.Debugf("add index/field/view/fragment to view.fragments: %s/%s/%s/%d", v.index, v.field, v.name, shard)
v.holder.Logger.Debugf("add index/field/view/fragment to view.fragments: %s/%s/%s/%d", v.index, v.field, v.name, shard)
mu.Lock()
v.fragments[frag.shard] = frag
v.addKnownShard(frag.shard)
@ -342,7 +342,7 @@ func (v *view) notifyIfNewShard(shard uint64) {
// Broadcast a message that a new max shard was just created.
err := v.broadcaster.SendSync(msg)
if err != nil {
v.logger.Printf("broadcasting create shard: %v", err)
v.holder.Logger.Printf("broadcasting create shard: %v", err)
}
close(broadcastChan)
}()
@ -352,19 +352,15 @@ func (v *view) notifyIfNewShard(shard uint64) {
select {
case <-broadcastChan:
case <-time.After(50 * time.Millisecond):
v.logger.Debugf("broadcasting create shard took >50ms")
v.holder.Logger.Debugf("broadcasting create shard took >50ms")
}
}
func (v *view) newFragment(path string, shard uint64) *fragment {
frag := newFragment(path, v.index, v.field, v.name, shard, v.flags())
frag := newFragment(v.holder, path, v.index, v.field, v.name, shard, v.flags())
frag.CacheType = v.cacheType
frag.CacheSize = v.cacheSize
frag.Logger = v.logger
frag.stats = v.stats
if v.snapshotQueue != nil {
frag.snapshotQueue = v.snapshotQueue
}
if v.fieldType == FieldTypeMutex {
frag.mutexVector = newRowsVector(frag)
} else if v.fieldType == FieldTypeBool {
@ -382,7 +378,7 @@ func (v *view) deleteFragment(shard uint64) error {
return ErrFragmentNotFound
}
v.logger.Printf("delete fragment: (%s/%s/%s) %d", v.index, v.field, v.name, shard)
v.holder.Logger.Printf("delete fragment: (%s/%s/%s) %d", v.index, v.field, v.name, shard)
// Close data files before deletion.
if err := fragment.Close(); err != nil {
@ -396,7 +392,7 @@ func (v *view) deleteFragment(shard uint64) error {
// Delete fragment cache file.
if err := os.Remove(fragment.cachePath()); err != nil {
v.logger.Printf("no cache file to delete for shard %d", shard)
v.holder.Logger.Printf("no cache file to delete for shard %d", shard)
}
delete(v.fragments, shard)

View file

@ -34,7 +34,7 @@ func mustOpenView(index, field, name string) *view {
CacheSize: DefaultCacheSize,
}
v := newView(path, index, field, name, fo)
v := newView(NewHolder(DefaultPartitionN), path, index, field, name, fo)
if err := v.open(); err != nil {
panic(err)
}