From 38de65eac032372e337cd950f6596cf377a70b7f Mon Sep 17 00:00:00 2001 From: Todd Gruben Date: Thu, 31 Jan 2019 14:09:16 -0600 Subject: [PATCH] only load shards that are applicable to node --- field.go | 5 +++++ holder.go | 8 ++++++++ index.go | 5 +++++ server.go | 3 +++ view.go | 16 ++++++++++++---- 5 files changed, 33 insertions(+), 4 deletions(-) diff --git a/field.go b/field.go index 38ba8ee76..b1384a71f 100644 --- a/field.go +++ b/field.go @@ -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 } diff --git a/holder.go b/holder.go index 967664168..3de4dfeb1 100644 --- a/holder.go +++ b/holder.go @@ -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 } diff --git a/index.go b/index.go index 696e81769..47ffe2cb3 100644 --- a/index.go +++ b/index.go @@ -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 } diff --git a/server.go b/server.go index 4902c87d8..22f9f390b 100644 --- a/server.go +++ b/server.go @@ -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 { diff --git a/view.go b/view.go index 658261037..7f10c54b8 100644 --- a/view.go +++ b/view.go @@ -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 {