mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-09-07 17:15:56 +00:00
Merge pull request #1919 from molecula/fb1163-etcd-source-of-truth
make etcd schema primary source of truth for indexes and fields
This commit is contained in:
commit
956c37ec85
4 changed files with 27 additions and 52 deletions
|
|
@ -312,6 +312,14 @@ type inMemSchemator struct {
|
|||
schema Schema
|
||||
}
|
||||
|
||||
// NewInMemSchemator instantiates an InMemSchemator
|
||||
// this allows new holders to have thier own, and not rely on a shared instance
|
||||
func NewInMemSchemator() *inMemSchemator {
|
||||
return &inMemSchemator{
|
||||
schema: make(Schema),
|
||||
}
|
||||
}
|
||||
|
||||
// Schema is an in-memory implementation of the Schemator Schema method.
|
||||
func (s *inMemSchemator) Schema(ctx context.Context) (Schema, error) {
|
||||
s.mu.RLock()
|
||||
|
|
|
|||
27
holder.go
27
holder.go
|
|
@ -9,7 +9,6 @@ import (
|
|||
"path/filepath"
|
||||
"runtime"
|
||||
"sort"
|
||||
"strings"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
|
|
@ -219,7 +218,7 @@ func DefaultHolderConfig() *HolderConfig {
|
|||
OpenIDAllocator: func(string, bool) (*idAllocator, error) { return &idAllocator{}, nil },
|
||||
TranslationSyncer: NopTranslationSyncer,
|
||||
Serializer: GobSerializer,
|
||||
Schemator: disco.InMemSchemator,
|
||||
Schemator: disco.NewInMemSchemator(),
|
||||
Sharder: disco.InMemSharder,
|
||||
CacheFlushInterval: defaultCacheFlushInterval,
|
||||
StatsClient: stats.NopStatsClient,
|
||||
|
|
@ -335,23 +334,7 @@ func (h *Holder) Open() error {
|
|||
}
|
||||
defer f.Close()
|
||||
|
||||
fis, err := f.Readdir(0)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "reading directory")
|
||||
}
|
||||
|
||||
for _, fi := range fis {
|
||||
// Skip files or hidden directories.
|
||||
if !fi.IsDir() || strings.HasPrefix(fi.Name(), ".") {
|
||||
continue
|
||||
}
|
||||
|
||||
// Only continue with indexes which are present in schema.
|
||||
idx, ok := schema[fi.Name()]
|
||||
if !ok {
|
||||
continue
|
||||
}
|
||||
|
||||
for idxKey, idx := range schema {
|
||||
// decode the CreateIndexMessage from the schema data in order to
|
||||
// get its metadata, such as CreateAt.
|
||||
cim, err := decodeCreateIndexMessage(h.serializer, idx.Data)
|
||||
|
|
@ -359,11 +342,11 @@ func (h *Holder) Open() error {
|
|||
return errors.Wrap(err, "decoding create index message")
|
||||
}
|
||||
|
||||
h.Logger.Printf("opening index: %s", filepath.Base(fi.Name()))
|
||||
h.Logger.Printf("opening index: %s", idxKey)
|
||||
|
||||
index, err := h.newIndex(h.IndexPath(filepath.Base(fi.Name())), filepath.Base(fi.Name()))
|
||||
index, err := h.newIndex(h.IndexPath(idxKey), idxKey)
|
||||
if errors.Cause(err) == ErrName {
|
||||
h.Logger.Errorf("opening index: %s, err=%s", fi.Name(), err)
|
||||
h.Logger.Errorf("opening index: %s, err=%s", idxKey, err)
|
||||
continue
|
||||
} else if err != nil {
|
||||
return errors.Wrap(err, "opening index")
|
||||
|
|
|
|||
|
|
@ -10,7 +10,7 @@ func mustHolderConfig() *HolderConfig {
|
|||
cfg := DefaultHolderConfig()
|
||||
cfg.StorageConfig.FsyncEnabled = false
|
||||
cfg.RBFConfig.FsyncEnabled = false
|
||||
cfg.Schemator = disco.InMemSchemator
|
||||
cfg.Schemator = disco.NewInMemSchemator()
|
||||
cfg.Sharder = disco.InMemSharder
|
||||
return cfg
|
||||
}
|
||||
|
|
|
|||
42
index.go
42
index.go
|
|
@ -77,7 +77,7 @@ func NewIndex(holder *Holder, path, name string) (*Index, error) {
|
|||
holder: holder,
|
||||
trackExistence: true,
|
||||
|
||||
Schemator: disco.InMemSchemator,
|
||||
Schemator: disco.NewInMemSchemator(),
|
||||
serializer: NopSerializer,
|
||||
|
||||
translateStores: make(map[int]TranslateStore),
|
||||
|
|
@ -264,43 +264,26 @@ func (i *Index) openFields(idx *disco.Index) error {
|
|||
}
|
||||
defer f.Close()
|
||||
|
||||
fis, err := f.Readdir(0)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "reading directory")
|
||||
}
|
||||
eg, ctx := errgroup.WithContext(context.Background())
|
||||
var mu sync.Mutex
|
||||
|
||||
if idx == nil {
|
||||
return nil
|
||||
}
|
||||
fileLoop:
|
||||
for _, loopFi := range fis {
|
||||
for fname, fld := range idx.Fields {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
break fileLoop
|
||||
default:
|
||||
fi := loopFi
|
||||
if !fi.IsDir() {
|
||||
continue
|
||||
}
|
||||
|
||||
var cfm *CreateFieldMessage = &CreateFieldMessage{}
|
||||
var err error
|
||||
|
||||
// Only continue with fields which are present in the provided,
|
||||
// non-nil index schema. The reason we have to check for idx != nil
|
||||
// here is because there are tests which call index.Open without
|
||||
// having a disco.Index available.
|
||||
if idx != nil {
|
||||
fld, ok := idx.Fields[fi.Name()]
|
||||
if !ok {
|
||||
continue
|
||||
}
|
||||
|
||||
// Decode the CreateFieldMessage from the schema data in order to
|
||||
// get its metadata.
|
||||
cfm, err = decodeCreateFieldMessage(i.holder.serializer, fld.Data)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "decoding create field message")
|
||||
}
|
||||
// Decode the CreateFieldMessage from the schema data in order to
|
||||
// get its metadata.
|
||||
cfm, err = decodeCreateFieldMessage(i.holder.serializer, fld.Data)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "decoding create field message")
|
||||
}
|
||||
|
||||
indexQueue <- struct{}{}
|
||||
|
|
@ -308,9 +291,9 @@ fileLoop:
|
|||
defer func() {
|
||||
<-indexQueue
|
||||
}()
|
||||
i.holder.Logger.Debugf("open field: %s", fi.Name())
|
||||
i.holder.Logger.Debugf("open field: %s", fname)
|
||||
|
||||
_, err := i.openField(&mu, cfm, fi.Name())
|
||||
_, err := i.openField(&mu, cfm, fname)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "opening field")
|
||||
}
|
||||
|
|
@ -319,6 +302,7 @@ fileLoop:
|
|||
})
|
||||
}
|
||||
}
|
||||
|
||||
err = eg.Wait()
|
||||
if err != nil {
|
||||
// Close any fields which got opened, since the overall
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue