From 8fe1b9a37ca46e9d163259a9b4810975e99ba251 Mon Sep 17 00:00:00 2001 From: Travis Date: Mon, 8 Feb 2021 23:19:26 -0600 Subject: [PATCH] store views in etcd via Schemator --- api.go | 2 +- cluster.go | 3 -- disco/disco.go | 25 +++++++++++++++- etcd/embed.go | 75 ++++++++++++++++++++++++++++++++++++++++++++--- field.go | 64 +++++++++++++++++++++++++++++++++------- holder.go | 79 +++++++++++++++++++++++++++++++++++++------------- index.go | 11 +++++-- pilosa.go | 2 ++ server.go | 20 ++++++------- 9 files changed, 227 insertions(+), 54 deletions(-) diff --git a/api.go b/api.go index 4906b80b8..c5a815c8b 100644 --- a/api.go +++ b/api.go @@ -1328,7 +1328,7 @@ func (api *API) Import(ctx context.Context, qcx *Qcx, req *ImportRequest, opts . return nil } -// Import bulk imports data into a particular index,field,shard. +// ImportWithTx bulk imports data into a particular index,field,shard. func (api *API) ImportWithTx(ctx context.Context, qcx *Qcx, req *ImportRequest, opts ...ImportOption) error { span, _ := tracing.StartSpanFromContext(ctx, "API.Import") defer span.Finish() diff --git a/cluster.go b/cluster.go index b4c891aa5..222e5c4c4 100644 --- a/cluster.go +++ b/cluster.go @@ -636,9 +636,6 @@ func (c *cluster) remoteSchema() (*Schema, error) { continue } - // TODO: replace following line by: - // ii, err := c.InternalClient.SchemaNode(context.Background(), &n.URI, true) - // after we ii, err := c.InternalClient.SchemaNode(context.Background(), &n.URI, true) if err != nil { return nil, errors.Wrapf(err, "getting schema from %s (%v)", n.ID, n.URI) diff --git a/disco/disco.go b/disco/disco.go index de77f88dd..0725d2d82 100644 --- a/disco/disco.go +++ b/disco/disco.go @@ -88,7 +88,14 @@ type Stator interface { // for each of its fields. type Index struct { Data []byte - Fields map[string][]byte + Fields map[string]*Field +} + +// Field is a struct which contains the data encoded for the field as well as +// for each of its views. +type Field struct { + Data []byte + Views map[string][]byte } type Schemator interface { @@ -99,6 +106,9 @@ type Schemator interface { Field(ctx context.Context, index, field string) ([]byte, error) CreateField(ctx context.Context, index, field string, val []byte) error DeleteField(ctx context.Context, index, field string) error + View(ctx context.Context, index, field, view string) ([]byte, error) + CreateView(ctx context.Context, index, field, view string, val []byte) error + DeleteView(ctx context.Context, index, field, view string) error } type Metadata interface { @@ -263,3 +273,16 @@ func (*nopSchemator) CreateField(ctx context.Context, index, field string, val [ // DeleteField is a no-op implementation of the Schemator DeleteField method. func (*nopSchemator) DeleteField(ctx context.Context, index, field string) error { return nil } + +// View is a no-op implementation of the Schemator View method. +func (*nopSchemator) View(ctx context.Context, index, field, view string) ([]byte, error) { + return nil, nil +} + +// CreateView is a no-op implementation of the Schemator CreateView method. +func (*nopSchemator) CreateView(ctx context.Context, index, field, view string, val []byte) error { + return nil +} + +// DeleteView is a no-op implementation of the Schemator DeleteView method. +func (*nopSchemator) DeleteView(ctx context.Context, index, field, view string) error { return nil } diff --git a/etcd/embed.go b/etcd/embed.go index 3caef98ca..0d2d919a4 100644 --- a/etcd/embed.go +++ b/etcd/embed.go @@ -484,25 +484,52 @@ func (e *Etcd) Schema(ctx context.Context) (map[string]*disco.Index, error) { return nil, err } + // The logic in the following for loop assumes that the list of keys is + // ordered such that index comes before field, which comes before view. + // For example: + // /index1 + // /index1/field1 + // /index1/field1/view1 + // /index1/field1/view2 + // /index1/field2 + // /index2 + // /index2/field1 + // m := make(map[string]*disco.Index) for i, k := range keys { tokens := strings.Split(strings.Trim(k, "/"), "/") // token[0] contains the schemaPrefix + + // token[1]: index index := tokens[1] if _, ok := m[index]; !ok { m[index] = &disco.Index{ Data: vals[i], - Fields: make(map[string][]byte), + Fields: make(map[string]*disco.Field), } + continue } flds := m[index].Fields + // token[2]: field if len(tokens) > 2 { field := tokens[2] - flds[field] = vals[i] + if _, ok := flds[field]; !ok { + flds[field] = &disco.Field{ + Data: vals[i], + Views: make(map[string][]byte), + } + continue + } + views := flds[field].Views + + // token[3]: view + if len(tokens) > 3 { + view := tokens[3] + views[view] = vals[i] + } } } - return m, nil } @@ -601,7 +628,7 @@ func (e *Etcd) Field(ctx context.Context, indexName string, name string) ([]byte func (e *Etcd) CreateField(ctx context.Context, indexName string, name string, val []byte) error { cli, err := e.client() if err != nil { - return errors.Wrap(err, "CreateIndex: creating client") + return errors.Wrap(err, "CreateField: creating client") } defer cli.Close() @@ -646,6 +673,46 @@ func (e *Etcd) DeleteField(ctx context.Context, indexname string, name string) e return errors.Wrap(err, "DeleteField") } +func (e *Etcd) View(ctx context.Context, indexName, fieldName, name string) ([]byte, error) { + key := schemaPrefix + indexName + "/" + fieldName + "/" + name + return e.getKeyBytes(ctx, key) +} + +// CreateView differs from CreateIndex and CreateField in that it does not +// return an error if the view already exists. If this logic needs to be +// changed, we likely need to introduce an ErrViewExists variable and return +// that. I decided not to do that now because it would require importing the +// etcd package into the pilosa root package. The better way to do that may be +// to define those error types in the disco package instead. +func (e *Etcd) CreateView(ctx context.Context, indexName, fieldName, name string, val []byte) error { + cli, err := e.client() + if err != nil { + return errors.Wrap(err, "CreateView: creating client") + } + defer cli.Close() + + key := schemaPrefix + indexName + "/" + fieldName + "/" + name + + // Set up Op to write view value as bytes. + op := clientv3.OpPut(key, "") + op.WithValueBytes(val) + + // Check for key existence, and execute Op within a transaction. + _, err = cli.KV.Txn(ctx). + If(clientv3util.KeyMissing(key)). + Then(op). + Commit() + if err != nil { + return errors.Wrap(err, "executing transaction") + } + + return nil +} + +func (e *Etcd) DeleteView(ctx context.Context, indexName, fieldName, name string) error { + return e.delKey(ctx, schemaPrefix+indexName+"/"+fieldName+"/"+name, false) +} + func (e *Etcd) putKey(ctx context.Context, key, val string, opts ...clientv3.OpOption) error { cli, err := e.client() if err != nil { diff --git a/field.go b/field.go index e5936503d..37a38a249 100644 --- a/field.go +++ b/field.go @@ -31,6 +31,7 @@ import ( "time" "github.com/gogo/protobuf/proto" + "github.com/pilosa/pilosa/v2/disco" "github.com/pilosa/pilosa/v2/internal" "github.com/pilosa/pilosa/v2/pql" "github.com/pilosa/pilosa/v2/roaring" @@ -102,6 +103,8 @@ type Field struct { broadcaster broadcaster Stats stats.StatsClient + schemator disco.Schemator + serializer Serializer // Field options. options FieldOptions @@ -367,6 +370,8 @@ func newField(holder *Holder, path, index, name string, opts FieldOption) (*Fiel broadcaster: NopBroadcaster, Stats: stats.NopStatsClient, + schemator: disco.NopSchemator, + serializer: NopSerializer, options: *applyDefaultOptions(&fo), @@ -1147,19 +1152,22 @@ func (f *Field) recalculateCaches() { // createViewIfNotExists returns the named view, creating it if necessary. // Additionally, a CreateViewMessage is sent to the cluster. func (f *Field) createViewIfNotExists(name string) (*view, error) { - view, created, err := f.createViewIfNotExistsBase(name) + cvm := &CreateViewMessage{ + Index: f.index, + Field: f.name, + View: name, + } + + // call this base method to isolate the mu.Lock and ensure we aren't holding + // the lock while calling SendSync below. + view, created, err := f.createViewIfNotExistsBase(cvm) if err != nil { return nil, err } if created { // Broadcast view creation to the cluster. - err = f.broadcaster.SendSync( - &CreateViewMessage{ - Index: f.index, - Field: f.name, - View: name, - }) + err := f.broadcaster.SendSync(cvm) if err != nil { return nil, errors.Wrap(err, "sending CreateView message") } @@ -1169,15 +1177,26 @@ func (f *Field) createViewIfNotExists(name string) (*view, error) { } // createViewIfNotExistsBase returns the named view, creating it if necessary. -// The returned bool indicates whether the view was created or not. -func (f *Field) createViewIfNotExistsBase(name string) (*view, bool, error) { +// One purpose of isolating this method from createViewIfNotExists() is that we +// need to enforce the mu.Lock on everything in this method, but we can't be +// holding the lock when broadcasting the CreateViewMessage view +// broadcaster.SendSync(); calling that SendSync() while holding the lock can +// result in a deadlock waiting on the remote node to give up its lock obtained +// by performing the same action. The returned bool indicates whether the view +// was created or not. +func (f *Field) createViewIfNotExistsBase(cvm *CreateViewMessage) (*view, bool, error) { f.mu.Lock() defer f.mu.Unlock() - if view := f.viewMap[name]; view != nil { + // Create the view in etcd as the system of record. + if err := f.persistView(context.Background(), cvm); err != nil { + return nil, false, errors.Wrap(err, "persisting view") + } + + if view := f.viewMap[cvm.View]; view != nil { return view, false, nil } - view := f.newView(f.viewPath(name), name) + view := f.newView(f.viewPath(cvm.View), cvm.View) if err := view.openEmpty(); err != nil { return nil, false, errors.Wrap(err, "opening view") @@ -1216,6 +1235,11 @@ func (f *Field) deleteView(name string) error { delete(f.viewMap, name) + // Delete the view from etcd as the system of record. + if err := f.schemator.DeleteView(context.TODO(), f.index, f.name, name); err != nil { + return errors.Wrapf(err, "deleting view from etcd: %s/%s/%s", f.index, f.name, name) + } + return nil } @@ -2200,3 +2224,21 @@ func bitDepthInt64(v int64) uint { func FormatQualifiedFieldName(index, field string) string { return fmt.Sprintf("%s\x00%s\x00", index, field) } + +// persistView stores the view information in etcd. +func (f *Field) persistView(ctx context.Context, cvm *CreateViewMessage) error { + if cvm.Index == "" { + return ErrIndexRequired + } else if cvm.Field == "" { + return ErrFieldRequired + } else if cvm.View == "" { + return ErrViewRequired + } + + if b, err := f.serializer.Marshal(cvm); err != nil { + return errors.Wrap(err, "marshaling") + } else if err := f.schemator.CreateView(ctx, cvm.Index, cvm.Field, cvm.View, b); err != nil { + return errors.Wrapf(err, "writing field to disco: %s/%s/%s", cvm.Index, cvm.Field, cvm.View) + } + return nil +} diff --git a/holder.go b/holder.go index 47345a65a..9ecf673d8 100644 --- a/holder.go +++ b/holder.go @@ -878,7 +878,7 @@ func (h *Holder) schema(ctx context.Context, includeViews bool) ([]*IndexInfo, e return nil, errors.Wrapf(err, "getting schema via schemator") } - for indexName, index := range schema { + for _, index := range schema { cim, err := h.decodeCreateIndexMessage(index.Data) if err != nil { return nil, errors.Wrap(err, "decoding CreateIndexMessage") @@ -891,28 +891,28 @@ func (h *Holder) schema(ctx context.Context, includeViews bool) ([]*IndexInfo, e ShardWidth: ShardWidth, Fields: make([]*FieldInfo, 0, len(index.Fields)), } - for fieldName, fieldData := range index.Fields { - createFieldMessage, err := h.decodeCreateFieldMessage(fieldData) + for _, field := range index.Fields { + cfm, err := h.decodeCreateFieldMessage(field.Data) if err != nil { return nil, errors.Wrap(err, "decoding CreateFieldMessage") } - if fieldName == existenceFieldName { + if cfm.Field == existenceFieldName { continue } fi := &FieldInfo{ - Name: fieldName, - CreatedAt: createFieldMessage.CreatedAt, - Options: *createFieldMessage.Meta, + Name: cfm.Field, + CreatedAt: cfm.CreatedAt, + Options: *cfm.Meta, } if includeViews { - // Because views are not stored in etcd, we still rely on the - // local representation of views. - if localField := h.Field(indexName, fieldName); localField != nil { - for _, view := range localField.views() { - fi.Views = append(fi.Views, &ViewInfo{Name: view.name}) + for _, viewData := range field.Views { + cvm, err := h.decodeCreateViewMessage(viewData) + if err != nil { + return nil, errors.Wrap(err, "decoding CreateViewMessage") } - sort.Sort(viewInfoSlice(fi.Views)) + fi.Views = append(fi.Views, &ViewInfo{Name: cvm.View}) } + sort.Sort(viewInfoSlice(fi.Views)) } di.Fields = append(di.Fields, fi) } @@ -1040,16 +1040,28 @@ func (h *Holder) LoadIndex(name string) (*Index, error) { // LoadField creates a field based on the information stored in schemator. // An error is returned if the field already exists. func (h *Holder) LoadField(index, field string) (*Field, error) { - h.mu.Lock() - defer h.mu.Unlock() - // Ensure field doesn't already exist. if h.Field(index, field) != nil { return nil, newConflictError(ErrFieldExists) } + + h.mu.Lock() + defer h.mu.Unlock() + return h.loadField(index, field) } +// LoadView creates a view based on the information stored in schemator. Unlike +// index and field, it is not considered an error if the view already exists. +func (h *Holder) LoadView(index, field, view string) (*view, error) { + // If the view already exists, just return with it here. + if v := h.view(index, field, view); v != nil { + return v, nil + } + + return h.loadView(index, field, view) +} + // CreateIndexAndBroadcast creates an index locally, then broadcasts the // creation to other nodes so they can create locally as well. An error is // returned if the index already exists. @@ -1181,8 +1193,6 @@ func (h *Holder) loadIndex(indexName string) (*Index, error) { func (h *Holder) loadField(indexName, fieldName string) (*Field, error) { b, err := h.schemator.Field(context.TODO(), indexName, fieldName) if err != nil { - // TODO: we may need to wrap with ConflictError if the error type is - // ErrIndexExists. return nil, errors.Wrapf(err, "getting field: %s/%s", indexName, fieldName) } @@ -1192,12 +1202,33 @@ func (h *Holder) loadField(indexName, fieldName string) (*Field, error) { return nil, errors.Errorf("local index not found: %s", indexName) } - createFieldMessage, err := h.decodeCreateFieldMessage(b) + cfm, err := h.decodeCreateFieldMessage(b) if err != nil { return nil, errors.Wrap(err, "decoding CreateFieldMessage") } - return idx.createFieldIfNotExists(fieldName, createFieldMessage.Meta) + // TODO: can this take cfm? + return idx.createFieldIfNotExists(fieldName, cfm.Meta) +} + +func (h *Holder) loadView(indexName, fieldName, viewName string) (*view, error) { + b, err := h.schemator.View(context.TODO(), indexName, fieldName, viewName) + if err != nil { + return nil, errors.Wrapf(err, "getting view: %s/%s/%s", indexName, fieldName, viewName) + } + + // Get field. + fld := h.Field(indexName, fieldName) + if fld == nil { + return nil, errors.Errorf("local field not found: %s/%s", indexName, fieldName) + } + + cvm, err := h.decodeCreateViewMessage(b) + if err != nil { + return nil, errors.Wrap(err, "decoding CreateFieldMessage") + } + + return fld.createViewIfNotExists(cvm.View) } func (h *Holder) newIndex(path, name string) (*Index, error) { @@ -2231,3 +2262,11 @@ func (h *Holder) decodeCreateFieldMessage(b []byte) (*CreateFieldMessage, error) } return &cfm, nil } + +func (h *Holder) decodeCreateViewMessage(b []byte) (*CreateViewMessage, error) { + var cvm CreateViewMessage + if err := h.serializer.Unmarshal(b, &cvm); err != nil { + return nil, errors.Wrap(err, "unmarshaling") + } + return &cvm, nil +} diff --git a/index.go b/index.go index ac5afedc2..e510be7eb 100644 --- a/index.go +++ b/index.go @@ -102,6 +102,9 @@ func NewIndex(holder *Holder, path, name string) (*Index, error) { holder: holder, trackExistence: true, + schemator: disco.NopSchemator, + serializer: NopSerializer, + translateStores: make(map[int]TranslateStore), translationSyncer: NopTranslationSyncer, @@ -523,7 +526,7 @@ func (i *Index) CreateField(name string, opts ...FieldOption) (*Field, error) { // Create the field in etcd as the system of record. if err := i.persistField(context.Background(), cfm); err != nil { - return nil, errors.Wrap(err, "persisting index") + return nil, errors.Wrap(err, "persisting field") } return i.createField(cfm, false) @@ -548,7 +551,7 @@ func (i *Index) CreateFieldAndBroadcast(cfm *CreateFieldMessage) (*Field, error) // Create the field in etcd as the system of record. if err := i.persistField(context.Background(), cfm); err != nil { - return nil, errors.Wrap(err, "persisting index") + return nil, errors.Wrap(err, "persisting field") } return i.createField(cfm, true) @@ -588,7 +591,7 @@ func (i *Index) CreateFieldIfNotExists(name string, opts ...FieldOption) (*Field // persistent storage. In that case, this will return an "index exists" // error, which in that case should return the index. TODO: We may need // to allow for that in the future. - return nil, errors.Wrap(err, "persisting index") + return nil, errors.Wrap(err, "persisting field") } return i.createField(cfm, false) @@ -697,6 +700,8 @@ func (i *Index) newField(path, name string) (*Field, error) { f.idx = i f.Stats = i.Stats f.broadcaster = i.broadcaster + f.schemator = i.schemator + f.serializer = i.serializer f.rowAttrStore = i.newAttrStore(filepath.Join(f.path, ".data")) f.OpenTranslateStore = i.OpenTranslateStore return f, nil diff --git a/pilosa.go b/pilosa.go index c98f886e0..dbd4e6567 100644 --- a/pilosa.go +++ b/pilosa.go @@ -53,6 +53,8 @@ var ( ErrInvalidBetweenValue = errors.New("invalid value for between operation") ErrDecimalOutOfRange = errors.New("decimal value out of range") + ErrViewRequired = errors.New("view required") + ErrViewExists = errors.New("view already exists") ErrInvalidView = errors.New("invalid view") ErrInvalidCacheType = errors.New("invalid cache type") diff --git a/server.go b/server.go index 800d4bb8c..216444a0e 100644 --- a/server.go +++ b/server.go @@ -415,12 +415,14 @@ func NewServer(opts ...ServerOption) (*Server, error) { metricInterval: 0, diagnosticInterval: 0, - disCo: disco.NopDisCo, - stator: disco.NopStator, - metadator: disco.NopMetadator, - resizer: disco.NopResizer, - noder: topology.NewEmptyLocalNoder(), - sharder: disco.NopSharder, + disCo: disco.NopDisCo, + stator: disco.NopStator, + metadator: disco.NopMetadator, + resizer: disco.NopResizer, + noder: topology.NewEmptyLocalNoder(), + sharder: disco.NopSharder, + schemator: disco.NopSchemator, + serializer: NopSerializer, confirmDownRetries: defaultConfirmDownRetries, confirmDownSleep: defaultConfirmDownSleep, @@ -807,11 +809,7 @@ func (s *Server) receiveMessage(m Message) error { } case *CreateViewMessage: - f := s.holder.Field(obj.Index, obj.Field) - if f == nil { - return fmt.Errorf("local field not found: %s", obj.Field) - } - if _, _, err := f.createViewIfNotExistsBase(obj.View); err != nil { + if _, err := s.holder.LoadView(obj.Index, obj.Field, obj.View); err != nil { return err }