write shards per node

This commit is contained in:
Kuba Podgórski 2021-05-10 13:54:18 +02:00
parent 2517ee1bde
commit b36de3146a
4 changed files with 53 additions and 42 deletions

View file

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

View file

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

View file

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

View file

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