store views in etcd via Schemator

This commit is contained in:
Travis 2021-02-08 23:19:26 -06:00
parent 86f45e9e62
commit 8fe1b9a37c
No known key found for this signature in database
GPG key ID: 37080CC2042BA34E
9 changed files with 227 additions and 54 deletions

2
api.go
View file

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

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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