only load shards that are applicable to node

This commit is contained in:
Todd Gruben 2019-01-31 14:09:16 -06:00 committed by Matt Jaffee
parent 66d750bd53
commit 38de65eac0
No known key found for this signature in database
GPG key ID: 08A3DFFF987B11BF
5 changed files with 33 additions and 4 deletions

View file

@ -80,6 +80,7 @@ type Field struct {
// Shards with data on any node in the cluster, according to this node.
remoteAvailableShards *roaring.Bitmap
shardValidator func(uint64)bool
logger logger.Logger
}
@ -207,6 +208,9 @@ func newField(path, index, name string, opts FieldOption) (*Field, error) {
remoteAvailableShards: roaring.NewBitmap(),
shardValidator:func(uint64)bool{
return true
},
logger: logger.NopLogger,
}
return f, nil
@ -754,6 +758,7 @@ func (f *Field) newView(path, name string) *view {
view.rowAttrStore = f.rowAttrStore
view.stats = f.Stats.WithTags(fmt.Sprintf("view:%s", name))
view.broadcaster = f.broadcaster
view.shardValidator = f.shardValidator
return view
}

View file

@ -77,6 +77,8 @@ type Holder struct {
// The interval at which the cached row ids are persisted to disk.
cacheFlushInterval time.Duration
shardValidatorfn func(index string ,shard uint64)bool
Logger logger.Logger
}
@ -123,6 +125,9 @@ func NewHolder() *Holder {
NewAttrStore: newNopAttrStore,
cacheFlushInterval: defaultCacheFlushInterval,
shardValidatorfn: func(index string,shard uint64)bool{
return true //default
},
Logger: logger.NopLogger,
}
@ -425,6 +430,9 @@ func (h *Holder) newIndex(path, name string) (*Index, error) {
index.broadcaster = h.broadcaster
index.newAttrStore = h.NewAttrStore
index.columnAttrs = h.NewAttrStore(filepath.Join(index.path, ".data"))
index.shardValidator = func (shard uint64) bool{
return h.shardValidatorfn(name,shard)
}
return index, nil
}

View file

@ -52,6 +52,7 @@ type Index struct {
broadcaster broadcaster
Stats stats.StatsClient
shardValidator func(uint64)bool
logger logger.Logger
}
@ -74,6 +75,9 @@ func NewIndex(path, name string) (*Index, error) {
broadcaster: NopBroadcaster,
Stats: stats.NopStatsClient,
logger: logger.NopLogger,
shardValidator: func(uint64)bool{
return true
},
trackExistence: true,
}, nil
}
@ -403,6 +407,7 @@ func (i *Index) newField(path, name string) (*Field, error) {
f.Stats = i.Stats.WithTags(fmt.Sprintf("field:%s", name))
f.broadcaster = i.broadcaster
f.rowAttrStore = i.newAttrStore(filepath.Join(f.path, ".data"))
f.shardValidator = i.shardValidator
return f, nil
}

View file

@ -321,6 +321,9 @@ func NewServer(opts ...ServerOption) (*Server, error) {
s.cluster.broadcaster = s
s.cluster.maxWritesPerRequest = s.maxWritesPerRequest
s.holder.broadcaster = s
s.holder.shardValidatorfn = func(index string, shard uint64) bool {
return s.cluster.ownsShard(s.nodeID, index, shard)
}
err = s.cluster.setup()
if err != nil {

16
view.go
View file

@ -52,10 +52,11 @@ type view struct {
// Fragments by shard.
fragments map[uint64]*fragment
broadcaster broadcaster
stats stats.StatsClient
rowAttrStore AttrStore
logger logger.Logger
broadcaster broadcaster
stats stats.StatsClient
rowAttrStore AttrStore
logger logger.Logger
shardValidator func(uint64) bool
}
// newView returns a new instance of View.
@ -75,6 +76,9 @@ func newView(path, index, field, name string, fieldOptions FieldOptions) *view {
broadcaster: NopBroadcaster,
stats: stats.NopStatsClient,
logger: logger.NopLogger,
shardValidator: func(uint64) bool {
return true
},
}
}
@ -132,6 +136,10 @@ func (v *view) openFragments() error {
if err != nil {
continue
}
//skip shard if not owned
if !v.shardValidator(shard) {
continue
}
frag := v.newFragment(v.fragmentPath(shard), shard)
if err := frag.Open(); err != nil {