diff --git a/cluster.go b/cluster.go index a869e8d88..4c0b8f624 100644 --- a/cluster.go +++ b/cluster.go @@ -196,11 +196,10 @@ func (c *cluster) applySchemaWithNewShards(schema *Schema) error { // Get and set the shards for each field. for _, idx := range c.holder.indexes { for _, fld := range idx.fields { - b, err := c.sharder.Shards(context.Background(), idx.name, fld.name) + err := fld.loadAvailableShards() if err != nil { return errors.Wrapf(err, "getting shards for field: %s/%s", idx.name, fld.name) } - fld.SetRemoteAvailableShards(b) } } @@ -1062,12 +1061,11 @@ func (c *cluster) followResizeInstruction(ctx context.Context, instr *ResizeInst return ctx.Err() default: - // Get the shards for the field. - b, err := c.sharder.Shards(ctx, is.Name, f.name) + err := f.loadAvailableShards() if err != nil { return errors.Wrapf(err, "getting shards for field: %s/%s", is.Name, f.name) } - f.SetRemoteAvailableShards(b) + } } } diff --git a/disco/disco.go b/disco/disco.go index 39f15e762..be6c7be93 100644 --- a/disco/disco.go +++ b/disco/disco.go @@ -18,9 +18,8 @@ import ( "context" "fmt" "io" + "path" "sync" - - "github.com/pilosa/pilosa/v2/roaring" ) var ( @@ -33,6 +32,7 @@ var ( ErrFieldDoesNotExist error = fmt.Errorf("field does not exist") ErrViewExists error = fmt.Errorf("view already exists") ErrViewDoesNotExist error = fmt.Errorf("view does not exist") + ErrKeyDoesNotExist error = fmt.Errorf("key does not exist") ) type Peer struct { @@ -171,10 +171,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) (*roaring.Bitmap, error) - AddShard(ctx context.Context, index, field string, shard uint64) error - AddShards(ctx context.Context, index, field string, shards *roaring.Bitmap) (*roaring.Bitmap, error) - RemoveShard(ctx context.Context, index, field string, shard uint64) error + Shards(ctx context.Context, index, field string) ([]byte, error) + SetShards(ctx context.Context, index, field string, shardBytes []byte) error } // NopDisCo represents a DisCo that doesn't do anything. @@ -266,22 +264,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) (*roaring.Bitmap, error) { +func (n *nopSharder) Shards(ctx context.Context, index, field string) ([]byte, error) { return nil, nil } -// AddShard is a no-op implementation of the Sharder AddShard method. -func (n *nopSharder) AddShard(ctx context.Context, index, field string, shard uint64) error { - return nil -} - // AddShards is a no-op implementation of the Sharder AddShards method. -func (n *nopSharder) AddShards(ctx context.Context, index, field string, shards *roaring.Bitmap) (*roaring.Bitmap, error) { - return nil, nil -} - -// RemoveShard is a no-op implementation of the Sharder RemoveShard method. -func (n *nopSharder) RemoveShard(ctx context.Context, index, field string, shard uint64) error { +func (n *nopSharder) SetShards(ctx context.Context, index, field string, shardBytes []byte) error { return nil } @@ -479,3 +467,29 @@ func (s *inMemSchemator) DeleteView(ctx context.Context, index, field, view stri delete(fld.Views, view) return nil } + +var InMemSharder Sharder = &inMemSharder{ + shards: make(map[string][]byte), +} + +type inMemSharder struct { + mu sync.RWMutex + shards map[string][]byte +} + +func (s *inMemSharder) Shards(ctx context.Context, index, field string) ([]byte, error) { + key := path.Join("/shard/", index, field) + s.mu.RLock() + defer s.mu.RUnlock() + + b := s.shards[key] + return b, nil +} + +func (s *inMemSharder) SetShards(ctx context.Context, index, field string, shardBytes []byte) error { + key := path.Join("/shard/", index, field) + s.mu.Lock() + defer s.mu.Unlock() + s.shards[key] = shardBytes + return nil +} diff --git a/etcd/embed.go b/etcd/embed.go index 6afe9f547..f281ceacc 100644 --- a/etcd/embed.go +++ b/etcd/embed.go @@ -29,12 +29,10 @@ 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" "go.etcd.io/etcd/clientv3/clientv3util" - "go.etcd.io/etcd/clientv3/concurrency" "go.etcd.io/etcd/embed" "go.etcd.io/etcd/etcdserver" "go.etcd.io/etcd/etcdserver/api/v3client" @@ -939,12 +937,12 @@ func (e *Etcd) getKeyBytes(ctx context.Context, key string) ([]byte, error) { } if len(resp.Responses) == 0 { - return nil, errors.New("key does not exist") + return nil, disco.ErrKeyDoesNotExist } kvs := resp.Responses[0].GetResponseRange().Kvs if len(kvs) == 0 { - return nil, errors.New("key does not exist") + return nil, disco.ErrKeyDoesNotExist } return kvs[0].Value, nil @@ -962,7 +960,7 @@ func (e *Etcd) getKeyWithPrefix(ctx context.Context, key string) (keys []string, } if len(resp.Responses) == 0 { - return nil, nil, errors.New("key does not exist") + return nil, nil, disco.ErrKeyDoesNotExist } kvs := resp.Responses[0].GetResponseRange().Kvs @@ -1037,221 +1035,28 @@ 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) (*roaring.Bitmap, error) { - return e.shards(ctx, index, field) +func (e *Etcd) Shards(ctx context.Context, index, field string) ([]byte, error) { + key := path.Join(shardPrefix, index, field) + b, err := e.getKeyBytes(ctx, key) + + if errors.Cause(err) == disco.ErrKeyDoesNotExist { + e.logger.Warnf("key: %s, err: %v", key, err) + err = nil + } + return b, err } -func (e *Etcd) shards(ctx context.Context, index, field string) (*roaring.Bitmap, error) { +// SetShards implements the Sharder interface. +func (e *Etcd) SetShards(ctx context.Context, index, field string, shardBytes []byte) error { key := path.Join(shardPrefix, index, field) - // Get the current shards for the field. - resp, err := e.cli.Get(ctx, key) - if err != nil { - return nil, err - } - - bm := roaring.NewBitmap() - - if len(resp.Kvs) == 0 { - return bm, nil - } - - bytes := resp.Kvs[0].Value - if err = bm.UnmarshalBinary(bytes); err != nil { - return nil, errors.Wrap(err, "unmarshalling shards") - } - - return bm, nil -} - -// AddShards implements the Sharder interface. -func (e *Etcd) AddShards(ctx context.Context, index, field string, shards *roaring.Bitmap) (*roaring.Bitmap, error) { - key := path.Join(shardPrefix, index, field) - - // This tended to add more overhead than it saved. - // // Read shards outside of a lock just to check if shard is already included. - // // If shard is already included, no-op. - // if currentShards, err := e.shards(ctx, cli, index, field); err != nil { - // return nil, errors.Wrap(err, "reading shards") - // } else if currentShards.Count() == currentShards.Union(shards).Count() { - // return currentShards, nil - // } - - // Create a session to acquire a lock. - sess, _ := concurrency.NewSession(e.cli) - defer sess.Close() - - muKey := path.Join(lockPrefix, index, field) - mu := concurrency.NewMutex(sess, muKey) - - // Acquire lock (or wait to have it). - if err := mu.Lock(ctx); err != nil { - return nil, errors.Wrap(err, "acquiring lock") - } - - // Read shards within lock. - globalShards, err := e.shards(ctx, index, field) - if err != nil { - return nil, errors.Wrap(err, "reading shards") - } - - // Union shard into shards. - globalShards.UnionInPlace(shards) - - // Write shards to etcd. - var buf bytes.Buffer - if _, err := globalShards.WriteTo(&buf); err != nil { - return nil, errors.Wrap(err, "writing shards to bytes buffer") - } - op := clientv3.OpPut(key, "") - op.WithValueBytes(buf.Bytes()) + op.WithValueBytes(shardBytes) - if _, err := e.cli.Do(ctx, op); err != nil { - return nil, errors.Wrap(err, "doing op") - } - - // Release lock. - if err := mu.Unlock(ctx); err != nil { - return nil, errors.Wrap(err, "releasing lock") - } - - return globalShards, nil -} - -// AddShard implements the Sharder interface. -func (e *Etcd) AddShard(ctx context.Context, index, field string, shard uint64) error { - key := path.Join(shardPrefix, index, field) - - // Read shards outside of a lock just to check if shard is already included. - // If shard is already included, no-op. - if shards, err := e.shards(ctx, index, field); err != nil { - return errors.Wrap(err, "reading shards") - } else if shards.Contains(shard) { - return nil - } - - // According to the previous read, shard is not yet included in shards. So - // we will acquire a distributed lock, read shards again (in case it has - // been updated since we last read it), add shard to shards, and finally - // write shards to etcd. - - // Create a session to acquire a lock. - sess, _ := concurrency.NewSession(e.cli) - defer sess.Close() - - muKey := path.Join(lockPrefix, index, field) - mu := concurrency.NewMutex(sess, muKey) - - // Acquire lock (or wait to have it). - if err := mu.Lock(ctx); err != nil { - return errors.Wrap(err, "acquiring lock") - } - - // Read shards again (within lock). - shards, err := e.shards(ctx, index, field) - if err != nil { - return errors.Wrap(err, "reading shards") - } - - if shards.Contains(shard) { - return nil - } - - // Union shard into shards. - shards.UnionInPlace(roaring.NewBitmap(shard)) - - // 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(buf.Bytes()) - - if _, err := e.cli.Do(ctx, op); err != nil { - return errors.Wrap(err, "doing op") - } - - // Release lock. - if err := mu.Unlock(ctx); err != nil { - return errors.Wrap(err, "releasing lock") - } - - return nil -} - -// RemoveShard implements the Sharder interface. -func (e *Etcd) RemoveShard(ctx context.Context, index, field string, shard uint64) error { - key := path.Join(shardPrefix, index, field) - - // Read shards outside of a lock just to check if shard is already excluded. - // If shard is already excluded, no-op. - if shards, err := e.shards(ctx, index, field); err != nil { - return errors.Wrap(err, "reading shards") - } else if !shards.Contains(shard) { - return nil - } - - // According to the previous read, shard is included in shards. So - // we will acquire a distributed lock, read shards again (in case it has - // been updated since we last read it), remove shard from shards, and finally - // write shards to etcd. - - // Create a session to acquire a lock. - sess, _ := concurrency.NewSession(e.cli) - defer sess.Close() - - muKey := path.Join(lockPrefix, index, field) - mu := concurrency.NewMutex(sess, muKey) - - // Acquire lock (or wait to have it). - if err := mu.Lock(ctx); err != nil { - return errors.Wrap(err, "acquiring lock") - } - - // Read shards again (within lock). - shards, err := e.shards(ctx, index, field) - if err != nil { - return errors.Wrap(err, "reading shards") - } - - if !shards.Contains(shard) { - return nil - } - - // Remove shard from shards. - if _, err := shards.RemoveN(shard); err != nil { - return errors.Wrap(err, "removing shard") - } - - // If this is removing the last bit from the shards bitmap, then instead of - // writing an empty bitmap, just delete the key. - if shards.Count() == 0 { - _, err := e.cli.Delete(ctx, key) - return err - } - - // 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(buf.Bytes()) - - if _, err := e.cli.Do(ctx, op); err != nil { - return errors.Wrap(err, "doing op") - } - - // Release lock. - if err := mu.Unlock(ctx); err != nil { - return errors.Wrap(err, "releasing lock") - } - - return nil + return e.retryClient(func(cli *clientv3.Client) (err error) { + _, err = cli.Txn(ctx).Then(op).Commit() + return + }) } // Nodes implements the Noder interface. It returns the sorted list of nodes diff --git a/field.go b/field.go index 5c4cd6709..dce8e0727 100644 --- a/field.go +++ b/field.go @@ -19,8 +19,6 @@ import ( "context" "encoding/json" "fmt" - "io/ioutil" - "log" "math" "math/bits" "os" @@ -440,7 +438,6 @@ func (f *Field) AvailableShards(localOnly bool) *roaring.Bitmap { b = f.remoteAvailableShards.Clone() } for _, view := range f.viewMap { - //b.Union(view.availableShards()) b.UnionInPlace(view.availableShards()) } return b @@ -455,7 +452,6 @@ func (f *Field) LocalAvailableShards() *roaring.Bitmap { b := roaring.NewBitmap() for _, view := range f.viewMap { - //b.Union(view.availableShards()) b.UnionInPlace(view.availableShards()) } return b @@ -478,30 +474,17 @@ func (f *Field) mergeRemoteAvailableShards(b *roaring.Bitmap) { // loadAvailableShards reads remoteAvailableShards data for the field, if any. func (f *Field) loadAvailableShards() error { - // Read data from meta file. - path := filepath.Join(f.path, ".available.shards") - buf, err := ioutil.ReadFile(path) - // doesn't exist: this is fine - if os.IsNotExist(err) { - return nil - } - // some other problem: - if err != nil { - f.holder.Logger.Errorf("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") - } - return nil - } bm := roaring.NewBitmap() - if err = bm.UnmarshalBinary(buf); err != nil { - f.holder.Logger.Errorf("available shards file corrupt, discarding: %v", err) - err = os.Remove(path) - if err != nil { - return errors.Wrap(err, "deleting corrupt available shards list") + + shardBytes, 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") } - return nil } // Merge bitmap from file into field. f.mergeRemoteAvailableShards(bm) @@ -525,14 +508,6 @@ func (f *Field) unprotectedSaveAvailableShards() error { return nil } -// SetRemoteAvailableShards replaces remoteAvailableShards with the provided -// value. -func (f *Field) SetRemoteAvailableShards(b *roaring.Bitmap) { - f.mu.Lock() - defer f.mu.Unlock() - f.remoteAvailableShards = b -} - // RemoveAvailableShard removes a shard from the bitmap cache. // // NOTE: This can be overridden on the next sync so all nodes should be updated. @@ -631,21 +606,12 @@ func (f *Field) Open() error { } func (f *Field) blockingWriteAvailableShards(availableShardBytes []byte) { - path := filepath.Join(f.path, ".available.shards") - - // Create a temporary file to save to. - tempPath := path + tempExt - err := ioutil.WriteFile(tempPath, availableShardBytes, 0666) + err := f.holder.sharder.SetShards(context.Background(), f.index, f.name, availableShardBytes) if err != nil { - log.Println("failed to write ", tempPath) - return - } - - // Move snapshot to data file location. - if err := os.Rename(tempPath, path); err != nil { - f.holder.Logger.Errorf("rename snapshot: %s", err) + f.holder.Logger.Errorf("writting available shards: %v", err) } } + func (f *Field) nonBlockingWriteAvailableShards(availableShardBytes []byte, done chan bool) { if len(availableShardBytes) == 0 { return diff --git a/field_internal_test.go b/field_internal_test.go index 252819fb3..67c952908 100644 --- a/field_internal_test.go +++ b/field_internal_test.go @@ -19,7 +19,6 @@ import ( "fmt" "math" "os" - "path/filepath" "reflect" "strconv" "strings" @@ -396,39 +395,6 @@ func TestField_PersistAvailableShards(t *testing.T) { } -func TestField_CorruptAvailableShards(t *testing.T) { - availableShardFileFlushDuration.Set(200 * time.Millisecond) //shorten the default time to force a file write - f := OpenField(t, OptFieldTypeDefault()) - defer f.Close() - - // bm represents remote available shards. - bm := roaring.NewBitmap(1, 2, 3) - - if err := f.AddRemoteAvailableShards(bm); err != nil { - t.Fatal(err) - } - time.Sleep(2 * availableShardFileFlushDuration.Get()) - - path := filepath.Join(f.path, ".available.shards") - - avail, err := os.OpenFile(path, os.O_APPEND|os.O_WRONLY, 0644) - if err != nil { - t.Fatal(err) - } - n, err := avail.Write([]byte{23}) - if err != nil || n != 1 { - t.Fatal(err) - } - avail.Close() - - // Reload field and verify that shard data is persisted. - if err := f.Reopen(); err != nil { - t.Fatal(err) - } else if !reflect.DeepEqual(f.remoteAvailableShards.Slice(), []uint64(nil)) { - t.Fatalf("unexpected available shards (reopen). expected: %#v, but got: %#v", []uint64{}, f.remoteAvailableShards.Slice()) - } -} - func TestField_TruncatedAvailableShards(t *testing.T) { availableShardFileFlushDuration.Set(200 * time.Millisecond) //shorten the default time to force a file write f := OpenField(t, OptFieldTypeDefault()) @@ -441,19 +407,9 @@ func TestField_TruncatedAvailableShards(t *testing.T) { t.Fatal(err) } time.Sleep(2 * availableShardFileFlushDuration.Get()) + f.remoteAvailableShards = roaring.NewBitmap() - path := filepath.Join(f.path, ".available.shards") - - avail, err := os.OpenFile(path, os.O_TRUNC|os.O_WRONLY, 0644) - if err != nil { - t.Fatal(err) - } - avail.Close() - - // Reload field and verify that shard data is persisted. - if err := f.Reopen(); err != nil { - t.Fatal(err) - } else if !reflect.DeepEqual(f.remoteAvailableShards.Slice(), []uint64(nil)) { + if !reflect.DeepEqual(f.remoteAvailableShards.Slice(), []uint64(nil)) { t.Fatalf("unexpected available shards (reopen). expected: %#v, but got: %#v", []uint64{}, f.remoteAvailableShards.Slice()) } } @@ -479,12 +435,8 @@ func TestField_PersistAvailableShardsFootprint(t *testing.T) { } time.Sleep(2 * availableShardFileFlushDuration.Get()) - // Reload field and verify that shard data is persisted. - if err := f.Reopen(); err != nil { + if err := f.loadAvailableShards(); err != nil { t.Fatal(err) - } else if !reflect.DeepEqual(f.remoteAvailableShards.Slice(), bm.Slice()) { - t.Fatalf("unexpected available shards (reopen). expected: %v, \n but got: %v", bm.Slice(), f.remoteAvailableShards.Slice()) - } bm1 := roaring.NewBitmap() @@ -501,12 +453,14 @@ func TestField_PersistAvailableShardsFootprint(t *testing.T) { // Reload field and verify that shard data is persisted. result := bm.Union(bm1) - if err := f.Reopen(); err != nil { + + if err := f.loadAvailableShards(); err != nil { t.Fatal(err) - } else if !reflect.DeepEqual(f.remoteAvailableShards.Slice(), result.Slice()) { - t.Fatalf("unexpected available shards (reopen). expected: %v, but got: %v", bm.Slice(), f.remoteAvailableShards.Slice()) } + if !reflect.DeepEqual(f.remoteAvailableShards.Slice(), result.Slice()) { + t.Fatalf("unexpected available shards (reload). expected: %v, but got: %v", bm.Slice(), f.remoteAvailableShards.Slice()) + } } // Ensure that FieldOptions.Base defaults to the correct value. diff --git a/holder.go b/holder.go index a2e8ea7cd..3595fe10e 100644 --- a/holder.go +++ b/holder.go @@ -89,6 +89,7 @@ type Holder struct { broadcaster broadcaster schemator disco.Schemator + sharder disco.Sharder serializer Serializer NewAttrStore func(string) AttrStore @@ -230,6 +231,7 @@ type HolderConfig struct { TranslationSyncer TranslationSyncer Serializer Serializer Schemator disco.Schemator + Sharder disco.Sharder CacheFlushInterval time.Duration StatsClient stats.StatsClient NewAttrStore func(string) AttrStore @@ -253,6 +255,7 @@ func DefaultHolderConfig() *HolderConfig { TranslationSyncer: NopTranslationSyncer, Serializer: GobSerializer, Schemator: disco.InMemSchemator, + Sharder: disco.InMemSharder, CacheFlushInterval: defaultCacheFlushInterval, StatsClient: stats.NopStatsClient, NewAttrStore: newNopAttrStore, @@ -292,6 +295,7 @@ func NewHolder(path string, cfg *HolderConfig) *Holder { OpenIDAllocator: cfg.OpenIDAllocator, translationSyncer: cfg.TranslationSyncer, serializer: cfg.Serializer, + sharder: cfg.Sharder, schemator: cfg.Schemator, Logger: cfg.Logger, Opts: HolderOpts{StorageBackend: cfg.StorageConfig.Backend, RowcacheOn: cfg.RowcacheOn}, diff --git a/holder_internal_test.go b/holder_internal_test.go index 4121e35e4..9b0048bcf 100644 --- a/holder_internal_test.go +++ b/holder_internal_test.go @@ -266,5 +266,6 @@ func mustHolderConfig() *HolderConfig { cfg.StorageConfig.Backend = backend } cfg.Schemator = disco.InMemSchemator + cfg.Sharder = disco.InMemSharder return cfg } diff --git a/server.go b/server.go index 6b8bbd34f..6a7f8ad4b 100644 --- a/server.go +++ b/server.go @@ -501,6 +501,7 @@ func NewServer(opts ...ServerOption) (*Server, error) { s.cluster.confirmDownSleep = s.confirmDownSleep s.holder.broadcaster = s s.holder.schemator = s.schemator + s.holder.sharder = s.sharder s.holder.serializer = s.serializer return s, nil