Write remote available shards to etcd, instead of local file.

This commit is contained in:
Kuba Podgórski 2021-05-07 15:33:18 +02:00
parent 1af85a818b
commit a79a36232f
8 changed files with 80 additions and 337 deletions

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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

View file

@ -266,5 +266,6 @@ func mustHolderConfig() *HolderConfig {
cfg.StorageConfig.Backend = backend
}
cfg.Schemator = disco.InMemSchemator
cfg.Sharder = disco.InMemSharder
return cfg
}

View file

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