diff --git a/disco/disco.go b/disco/disco.go index be6c7be93..2e20de9f9 100644 --- a/disco/disco.go +++ b/disco/disco.go @@ -20,6 +20,8 @@ import ( "io" "path" "sync" + + "github.com/pilosa/pilosa/v2/roaring" ) var ( @@ -171,8 +173,8 @@ type Resizer interface { // Sharder is an interface used to maintain the set of availableShards bitmaps // per field. type Sharder interface { - Shards(ctx context.Context, index, field string) ([]byte, error) - SetShards(ctx context.Context, index, field string, shardBytes []byte) error + Shards(ctx context.Context, index, field string) (*roaring.Bitmap, error) + SetShards(ctx context.Context, index, field string, shards *roaring.Bitmap) error } // NopDisCo represents a DisCo that doesn't do anything. @@ -264,12 +266,12 @@ var NopSharder Sharder = &nopSharder{} type nopSharder struct{} // Shards is a no-op implementation of the Sharder Shards method. -func (n *nopSharder) Shards(ctx context.Context, index, field string) ([]byte, error) { +func (n *nopSharder) Shards(ctx context.Context, index, field string) (*roaring.Bitmap, error) { return nil, nil } // AddShards is a no-op implementation of the Sharder AddShards method. -func (n *nopSharder) SetShards(ctx context.Context, index, field string, shardBytes []byte) error { +func (n *nopSharder) SetShards(ctx context.Context, index, field string, shards *roaring.Bitmap) error { return nil } @@ -469,27 +471,30 @@ func (s *inMemSchemator) DeleteView(ctx context.Context, index, field, view stri } var InMemSharder Sharder = &inMemSharder{ - shards: make(map[string][]byte), + shards: make(map[string]*roaring.Bitmap), } type inMemSharder struct { mu sync.RWMutex - shards map[string][]byte + shards map[string]*roaring.Bitmap } -func (s *inMemSharder) Shards(ctx context.Context, index, field string) ([]byte, error) { +func (s *inMemSharder) Shards(ctx context.Context, index, field string) (*roaring.Bitmap, error) { key := path.Join("/shard/", index, field) s.mu.RLock() defer s.mu.RUnlock() b := s.shards[key] + if b == nil { + return roaring.NewBitmap(), nil + } return b, nil } -func (s *inMemSharder) SetShards(ctx context.Context, index, field string, shardBytes []byte) error { +func (s *inMemSharder) SetShards(ctx context.Context, index, field string, shards *roaring.Bitmap) error { key := path.Join("/shard/", index, field) s.mu.Lock() defer s.mu.Unlock() - s.shards[key] = shardBytes + s.shards[key] = shards return nil } diff --git a/etcd/embed.go b/etcd/embed.go index 967d33411..63cbdc7b5 100644 --- a/etcd/embed.go +++ b/etcd/embed.go @@ -29,6 +29,7 @@ import ( "github.com/pilosa/pilosa/v2/disco" "github.com/pilosa/pilosa/v2/logger" + "github.com/pilosa/pilosa/v2/roaring" "github.com/pilosa/pilosa/v2/topology" "github.com/pkg/errors" "go.etcd.io/etcd/clientv3" @@ -1034,23 +1035,39 @@ func memberAdd(cli *clientv3.Client, peerURL string) (id uint64, name string) { } // Shards implements the Sharder interface. -func (e *Etcd) Shards(ctx context.Context, index, field string) ([]byte, error) { +func (e *Etcd) Shards(ctx context.Context, index, field string) (*roaring.Bitmap, error) { key := path.Join(shardPrefix, index, field) - b, err := e.getKeyBytes(ctx, key) + _, vals, err := e.getKeyWithPrefix(ctx, key) + bm := roaring.NewBitmap() if errors.Cause(err) == disco.ErrKeyDoesNotExist { e.logger.Warnf("key: %s, err: %v", key, err) - err = nil + return bm, nil } - return b, err + + for _, v := range vals { + b := roaring.NewBitmap() + if err = b.UnmarshalBinary(v); err != nil { + return nil, errors.Wrap(err, "unmarshalling shards") + } + bm.UnionInPlace(b) + } + + return bm, nil } // SetShards implements the Sharder interface. -func (e *Etcd) SetShards(ctx context.Context, index, field string, shardBytes []byte) error { - key := path.Join(shardPrefix, index, field) +func (e *Etcd) SetShards(ctx context.Context, index, field string, shards *roaring.Bitmap) error { + key := path.Join(shardPrefix, index, field, e.e.Server.ID().String()) + + // Write shards to etcd. + var buf bytes.Buffer + if _, err := shards.WriteTo(&buf); err != nil { + return errors.Wrap(err, "writing shards to bytes buffer") + } op := clientv3.OpPut(key, "") - op.WithValueBytes(shardBytes) + op.WithValueBytes(buf.Bytes()) return e.retryClient(func(cli *clientv3.Client) (err error) { _, err = cli.Txn(ctx).Then(op).Commit() diff --git a/field.go b/field.go index dce8e0727..fe154c78e 100644 --- a/field.go +++ b/field.go @@ -15,7 +15,6 @@ package pilosa import ( - "bytes" "context" "encoding/json" "fmt" @@ -127,7 +126,7 @@ type Field struct { // Synchronization primitives needed for async writing of // the remoteAvailableShards - availableShardChan chan []byte + availableShardChan chan *roaring.Bitmap doneChan chan struct{} wg sync.WaitGroup } @@ -474,20 +473,13 @@ func (f *Field) mergeRemoteAvailableShards(b *roaring.Bitmap) { // loadAvailableShards reads remoteAvailableShards data for the field, if any. func (f *Field) loadAvailableShards() error { - bm := roaring.NewBitmap() - - shardBytes, err := f.holder.sharder.Shards(context.Background(), f.index, f.name) + shards, err := f.holder.sharder.Shards(context.Background(), f.index, f.name) if err != nil { return errors.Wrap(err, "loading available shards") } - if shardBytes != nil { - if err = bm.UnmarshalBinary(shardBytes); err != nil { - return errors.Wrap(err, "available shards corrupt") - } - } // Merge bitmap from file into field. - f.mergeRemoteAvailableShards(bm) + f.mergeRemoteAvailableShards(shards) return nil } @@ -500,11 +492,8 @@ func (f *Field) saveAvailableShards() error { } func (f *Field) unprotectedSaveAvailableShards() error { - var buf bytes.Buffer - if _, err := f.remoteAvailableShards.WriteTo(&buf); err != nil { - return errors.Wrap(err, "rendering available shards ") - } - f.availableShardChan <- buf.Bytes() + f.remoteAvailableShards.Optimize() + f.availableShardChan <- f.remoteAvailableShards return nil } @@ -590,7 +579,7 @@ func (f *Field) Open() error { } } - f.availableShardChan = make(chan []byte) + f.availableShardChan = make(chan *roaring.Bitmap) f.doneChan = make(chan struct{}) f.wg.Add(1) go f.writeAvailableShards() @@ -605,19 +594,19 @@ func (f *Field) Open() error { return nil } -func (f *Field) blockingWriteAvailableShards(availableShardBytes []byte) { - err := f.holder.sharder.SetShards(context.Background(), f.index, f.name, availableShardBytes) +func (f *Field) blockingWriteAvailableShards(availableShards *roaring.Bitmap) { + err := f.holder.sharder.SetShards(context.Background(), f.index, f.name, availableShards) if err != nil { f.holder.Logger.Errorf("writting available shards: %v", err) } } -func (f *Field) nonBlockingWriteAvailableShards(availableShardBytes []byte, done chan bool) { - if len(availableShardBytes) == 0 { +func (f *Field) nonBlockingWriteAvailableShards(availableShards *roaring.Bitmap, done chan bool) { + if availableShards == nil { return } go func() { - f.blockingWriteAvailableShards(availableShardBytes) + f.blockingWriteAvailableShards(availableShards) done <- true }() } @@ -625,7 +614,7 @@ func (f *Field) nonBlockingWriteAvailableShards(availableShardBytes []byte, done func (f *Field) writeAvailableShards() { defer f.wg.Done() ticker := time.NewTicker(availableShardFileFlushDuration.Get()) - var data []byte + var data *roaring.Bitmap tracker := make(chan bool) writing := false @@ -634,7 +623,7 @@ func (f *Field) writeAvailableShards() { case newdata := <-f.availableShardChan: data = newdata case <-ticker.C: - if len(data) > 0 { + if data != nil { if !writing { writing = true f.nonBlockingWriteAvailableShards(data, tracker) @@ -647,7 +636,7 @@ func (f *Field) writeAvailableShards() { if writing { //wait to writing is complete <-tracker } - if len(data) > 0 { + if data != nil { f.blockingWriteAvailableShards(data) } alive = false diff --git a/field_test.go b/field_test.go index af58039a1..3deef292b 100644 --- a/field_test.go +++ b/field_test.go @@ -212,7 +212,7 @@ func TestField_AvailableShards(t *testing.T) { idx := test.MustOpenIndex(t) defer idx.Close() - f, err := idx.CreateField("f", pilosa.OptFieldTypeDefault()) + f, err := idx.CreateField("fld-shards", pilosa.OptFieldTypeDefault()) if err != nil { t.Fatal(err) }