mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-10-07 11:27:50 +00:00
refactoring stores back into index/field
This commit is contained in:
parent
7215bfd16c
commit
b3e86e8394
12 changed files with 342 additions and 221 deletions
104
api.go
104
api.go
|
|
@ -545,7 +545,7 @@ func (api *API) ExportCSV(ctx context.Context, indexName string, fieldName strin
|
|||
var err error
|
||||
|
||||
if field.keys() {
|
||||
if rowStr, err = api.cluster.translateFieldID(indexName, fieldName, rowID); err != nil {
|
||||
if rowStr, err = field.TranslateStore().TranslateID(rowID); err != nil {
|
||||
return errors.Wrap(err, "translating row")
|
||||
}
|
||||
} else {
|
||||
|
|
@ -553,7 +553,9 @@ func (api *API) ExportCSV(ctx context.Context, indexName string, fieldName strin
|
|||
}
|
||||
|
||||
if index.Keys() {
|
||||
if colStr, err = api.cluster.translateIndexPartitionID(ctx, indexName, api.cluster.idPartition(indexName, columnID), columnID); err != nil {
|
||||
if store := index.TranslateStore(api.cluster.idPartition(indexName, columnID)); store == nil {
|
||||
return errors.Wrap(err, "partition does not exist")
|
||||
} else if colStr, err = store.TranslateID(columnID); err != nil {
|
||||
return errors.Wrap(err, "translating column")
|
||||
}
|
||||
} else {
|
||||
|
|
@ -963,7 +965,7 @@ func (api *API) Import(ctx context.Context, req *ImportRequest, opts ...ImportOp
|
|||
if len(req.RowIDs) != 0 {
|
||||
return errors.New("row ids cannot be used because field uses string keys")
|
||||
}
|
||||
if req.RowIDs, err = api.cluster.translateFieldKeys(req.Index, req.Field, req.RowKeys); err != nil {
|
||||
if req.RowIDs, err = field.TranslateStore().TranslateKeys(req.RowKeys); err != nil {
|
||||
return errors.Wrap(err, "translating rows")
|
||||
}
|
||||
}
|
||||
|
|
@ -1368,10 +1370,67 @@ func (api *API) Info() serverInfo {
|
|||
}
|
||||
|
||||
// GetTranslateEntryReader provides an entry reader for key translation logs starting at offset.
|
||||
func (api *API) GetTranslateEntryReader(ctx context.Context, offsets TranslateOffsetMap) (TranslateEntryReader, error) {
|
||||
func (api *API) GetTranslateEntryReader(ctx context.Context, offsets TranslateOffsetMap) (_ TranslateEntryReader, err error) {
|
||||
span, ctx := tracing.StartSpanFromContext(ctx, "API.GetTranslateEntryReader")
|
||||
defer span.Finish()
|
||||
return api.cluster.translateEntryReader(ctx, offsets)
|
||||
|
||||
// Ensure all readers are cleaned up if any error.
|
||||
var a []TranslateEntryReader
|
||||
defer func() {
|
||||
if err != nil {
|
||||
for i := range a {
|
||||
a[i].Close()
|
||||
}
|
||||
}
|
||||
}()
|
||||
|
||||
// Fetch all index partition readers.
|
||||
for indexName, indexMap := range offsets {
|
||||
index := api.holder.Index(indexName)
|
||||
if index == nil {
|
||||
return nil, ErrIndexNotFound
|
||||
}
|
||||
|
||||
for partitionID, offset := range indexMap.Partitions {
|
||||
store := index.TranslateStore(partitionID)
|
||||
if store == nil {
|
||||
return nil, ErrTranslateStoreNotFound
|
||||
}
|
||||
|
||||
r, err := store.EntryReader(ctx, uint64(offset))
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "index partition translate reader")
|
||||
}
|
||||
a = append(a, r)
|
||||
}
|
||||
}
|
||||
|
||||
// Fetch all field readers.
|
||||
for indexName, indexMap := range offsets {
|
||||
index := api.holder.Index(indexName)
|
||||
if index == nil {
|
||||
return nil, ErrIndexNotFound
|
||||
}
|
||||
|
||||
for fieldName, offset := range indexMap.Fields {
|
||||
field := index.Field(fieldName)
|
||||
if field == nil {
|
||||
return nil, ErrIndexNotFound
|
||||
}
|
||||
|
||||
r, err := field.TranslateStore().EntryReader(ctx, uint64(offset))
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "field translate reader")
|
||||
}
|
||||
a = append(a, r)
|
||||
}
|
||||
}
|
||||
|
||||
return NewMultiTranslateEntryReader(ctx, a), nil
|
||||
}
|
||||
|
||||
func (api *API) TranslateIndexKey(ctx context.Context, indexName string, key string) (uint64, error) {
|
||||
return api.cluster.translateIndexKey(ctx, indexName, key)
|
||||
}
|
||||
|
||||
// TranslateKeys handles a TranslateKeyRequest.
|
||||
|
|
@ -1390,7 +1449,9 @@ func (api *API) TranslateKeys(ctx context.Context, r io.Reader) (_ []byte, err e
|
|||
return nil, err
|
||||
}
|
||||
} else {
|
||||
if ids, err = api.cluster.translateFieldKeys(req.Index, req.Field, req.Keys); err != nil {
|
||||
if field := api.holder.Field(req.Index, req.Field); field == nil {
|
||||
return nil, ErrFieldNotFound
|
||||
} else if ids, err = field.TranslateStore().TranslateKeys(req.Keys); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
}
|
||||
|
|
@ -1403,6 +1464,37 @@ func (api *API) TranslateKeys(ctx context.Context, r io.Reader) (_ []byte, err e
|
|||
return buf, nil
|
||||
}
|
||||
|
||||
// TranslateIDs handles a TranslateIDRequest.
|
||||
func (api *API) TranslateIDs(ctx context.Context, r io.Reader) (_ []byte, err error) {
|
||||
var req TranslateIDsRequest
|
||||
if buf, err := ioutil.ReadAll(r); err != nil {
|
||||
return nil, NewBadRequestError(errors.Wrap(err, "read translate ids request error"))
|
||||
} else if err := api.Serializer.Unmarshal(buf, &req); err != nil {
|
||||
return nil, NewBadRequestError(errors.Wrap(err, "unmarshal translate ids request error"))
|
||||
}
|
||||
|
||||
// Lookup store for either index or field and translate ids.
|
||||
var keys []string
|
||||
if req.Field == "" {
|
||||
if keys, err = api.cluster.translateIndexIDs(ctx, req.Index, req.IDs); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
} else {
|
||||
if field := api.holder.Field(req.Index, req.Field); field == nil {
|
||||
return nil, ErrFieldNotFound
|
||||
} else if keys, err = field.TranslateStore().TranslateIDs(req.IDs); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
}
|
||||
|
||||
// Encode response.
|
||||
buf, err := api.Serializer.Marshal(&TranslateIDsResponse{Keys: keys})
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "translate ids response encoding error")
|
||||
}
|
||||
return buf, nil
|
||||
}
|
||||
|
||||
// PrimaryReplicaNodeURL returns the URL of the cluster's primary replica.
|
||||
func (api *API) PrimaryReplicaNodeURL() url.URL {
|
||||
node := api.cluster.PrimaryReplicaNode()
|
||||
|
|
|
|||
341
cluster.go
341
cluster.go
|
|
@ -215,10 +215,6 @@ type cluster struct { // nolint: maligned
|
|||
holder *Holder
|
||||
broadcaster broadcaster
|
||||
|
||||
// translation stores
|
||||
indexTranslateStoreMap map[string]map[int]TranslateStore // store by index+partition
|
||||
fieldTranslateStoreMap map[string]map[string]TranslateStore // store by index+field
|
||||
|
||||
joiningLeavingNodes chan nodeAction
|
||||
|
||||
// joining is held open until this node
|
||||
|
|
@ -240,9 +236,7 @@ type cluster struct { // nolint: maligned
|
|||
|
||||
InternalClient InternalClient
|
||||
|
||||
// Instantiates new translation stores
|
||||
OpenTranslateStore OpenTranslateStoreFunc
|
||||
OpenTranslateReader OpenTranslateReaderFunc
|
||||
// OpenTranslateReader OpenTranslateReaderFunc
|
||||
}
|
||||
|
||||
// newCluster returns a new instance of Cluster with defaults.
|
||||
|
|
@ -257,13 +251,8 @@ func newCluster() *cluster {
|
|||
closing: make(chan struct{}),
|
||||
joining: make(chan struct{}),
|
||||
|
||||
indexTranslateStoreMap: make(map[string]map[int]TranslateStore),
|
||||
fieldTranslateStoreMap: make(map[string]map[string]TranslateStore),
|
||||
|
||||
InternalClient: newNopInternalClient(),
|
||||
|
||||
OpenTranslateStore: OpenInMemTranslateStore,
|
||||
|
||||
logger: logger.NopLogger,
|
||||
}
|
||||
}
|
||||
|
|
@ -1011,10 +1000,6 @@ func (c *cluster) setup() error {
|
|||
if err != nil {
|
||||
return errors.Wrap(err, "adding local node")
|
||||
}
|
||||
|
||||
if err := c.updateTranslateStores(); err != nil {
|
||||
return err
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
|
|
@ -2024,11 +2009,6 @@ func (c *cluster) mergeClusterStatus(cs *ClusterStatus) error {
|
|||
}
|
||||
}
|
||||
|
||||
// Open appropriate translate stores.
|
||||
if err := c.updateTranslateStores(); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
c.unprotectedSetState(cs.State)
|
||||
|
||||
c.markAsJoined()
|
||||
|
|
@ -2085,6 +2065,145 @@ func (c *cluster) setStatic(hosts []string) error {
|
|||
return nil
|
||||
}
|
||||
|
||||
func (c *cluster) translateIndexKey(ctx context.Context, indexName string, key string) (uint64, error) {
|
||||
keyMap, err := c.translateIndexKeySet(ctx, indexName, map[string]struct{}{key: struct{}{}})
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
return keyMap[key], nil
|
||||
}
|
||||
|
||||
func (c *cluster) translateIndexKeys(ctx context.Context, indexName string, keys []string) ([]uint64, error) {
|
||||
keySet := make(map[string]struct{})
|
||||
for _, key := range keys {
|
||||
keySet[key] = struct{}{}
|
||||
}
|
||||
|
||||
keyMap, err := c.translateIndexKeySet(ctx, indexName, keySet)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
ids := make([]uint64, len(keys))
|
||||
for i := range keys {
|
||||
ids[i] = keyMap[keys[i]]
|
||||
}
|
||||
return ids, nil
|
||||
}
|
||||
|
||||
func (c *cluster) translateIndexKeySet(ctx context.Context, indexName string, keySet map[string]struct{}) (map[string]uint64, error) {
|
||||
keyMap := make(map[string]uint64)
|
||||
|
||||
idx := c.holder.Index(indexName)
|
||||
if idx == nil {
|
||||
return nil, ErrIndexNotFound
|
||||
}
|
||||
|
||||
// Split keys by partition.
|
||||
keysByPartition := make(map[int][]string, c.partitionN)
|
||||
for key := range keySet {
|
||||
partitionID := c.keyPartition(indexName, key)
|
||||
keysByPartition[partitionID] = append(keysByPartition[partitionID], key)
|
||||
}
|
||||
|
||||
// Translate keys by partition.
|
||||
var g errgroup.Group
|
||||
var mu sync.Mutex
|
||||
for partitionID := range keysByPartition {
|
||||
keys := keysByPartition[partitionID]
|
||||
|
||||
g.Go(func() (err error) {
|
||||
var ids []uint64
|
||||
if c.ownsPartition(c.Node.ID, partitionID) {
|
||||
if ids, err = idx.TranslateStore(partitionID).TranslateKeys(keys); err != nil {
|
||||
return err
|
||||
}
|
||||
} else {
|
||||
nodes := c.partitionNodes(partitionID)
|
||||
if ids, err = c.InternalClient.TranslateKeysNode(ctx, &nodes[0].URI, indexName, "", keys); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
||||
mu.Lock()
|
||||
defer mu.Unlock()
|
||||
for i := range keys {
|
||||
keyMap[keys[i]] = ids[i]
|
||||
}
|
||||
return nil
|
||||
})
|
||||
}
|
||||
if err := g.Wait(); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return keyMap, nil
|
||||
}
|
||||
|
||||
func (c *cluster) translateIndexIDs(ctx context.Context, indexName string, ids []uint64) ([]string, error) {
|
||||
idSet := make(map[uint64]struct{})
|
||||
for _, id := range ids {
|
||||
idSet[id] = struct{}{}
|
||||
}
|
||||
|
||||
idMap, err := c.translateIndexIDSet(ctx, indexName, idSet)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
keys := make([]string, len(ids))
|
||||
for i := range ids {
|
||||
keys[i] = idMap[ids[i]]
|
||||
}
|
||||
return keys, nil
|
||||
}
|
||||
|
||||
func (c *cluster) translateIndexIDSet(ctx context.Context, indexName string, idSet map[uint64]struct{}) (map[uint64]string, error) {
|
||||
idMap := make(map[uint64]string)
|
||||
|
||||
index := c.holder.Index(indexName)
|
||||
if index == nil {
|
||||
return nil, ErrIndexNotFound
|
||||
}
|
||||
|
||||
// Split ids by partition.
|
||||
idsByPartition := make(map[int][]uint64, c.partitionN)
|
||||
for id := range idSet {
|
||||
partitionID := c.idPartition(indexName, id)
|
||||
idsByPartition[partitionID] = append(idsByPartition[partitionID], id)
|
||||
}
|
||||
|
||||
// Translate ids by partition.
|
||||
var g errgroup.Group
|
||||
var mu sync.Mutex
|
||||
for partitionID, ids := range idsByPartition {
|
||||
g.Go(func() (err error) {
|
||||
var keys []string
|
||||
if c.ownsPartition(c.Node.ID, partitionID) {
|
||||
if keys, err = index.TranslateStore(partitionID).TranslateIDs(ids); err != nil {
|
||||
return err
|
||||
}
|
||||
} else {
|
||||
nodes := c.partitionNodes(partitionID)
|
||||
if keys, err = c.InternalClient.TranslateIDsNode(ctx, &nodes[0].URI, indexName, "", ids); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
||||
mu.Lock()
|
||||
defer mu.Unlock()
|
||||
for i := range ids {
|
||||
idMap[ids[i]] = keys[i]
|
||||
}
|
||||
return nil
|
||||
})
|
||||
}
|
||||
if err := g.Wait(); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return idMap, nil
|
||||
}
|
||||
|
||||
/*
|
||||
// updateTranslateStores starts stores for partitions & fields owned by this
|
||||
// node and stops stores for ones this node does not own.
|
||||
func (c *cluster) updateTranslateStores() error {
|
||||
|
|
@ -2115,23 +2234,6 @@ func (c *cluster) closeUnownedTranslateStores() error {
|
|||
}
|
||||
}
|
||||
|
||||
// Close index partition stores.
|
||||
for indexName, m := range c.fieldTranslateStoreMap {
|
||||
idx := c.holder.Index(indexName)
|
||||
|
||||
for fieldName, store := range m {
|
||||
if idx != nil && idx.Field(fieldName) != nil {
|
||||
continue
|
||||
}
|
||||
if err := store.Close(); err != nil {
|
||||
return err
|
||||
}
|
||||
delete(m, fieldName)
|
||||
}
|
||||
if len(c.fieldTranslateStoreMap) == 0 {
|
||||
delete(c.fieldTranslateStoreMap, indexName)
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
|
|
@ -2157,26 +2259,6 @@ func (c *cluster) openOwnedTranslateStores() error {
|
|||
}
|
||||
}
|
||||
|
||||
// Open field stores.
|
||||
for _, index := range c.holder.Indexes() {
|
||||
m := c.fieldTranslateStoreMap[index.Name()]
|
||||
if m == nil {
|
||||
m = make(map[string]TranslateStore)
|
||||
c.fieldTranslateStoreMap[index.Name()] = m
|
||||
}
|
||||
|
||||
for _, field := range index.Fields() {
|
||||
if m[field.Name()] != nil {
|
||||
continue
|
||||
}
|
||||
store, err := c.OpenTranslateStore(field.TranslateStorePath(), index.Name(), field.Name(), 0)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
m[field.Name()] = store
|
||||
}
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
|
|
@ -2196,101 +2278,7 @@ func (c *cluster) fieldTranslateStore(index, field string) TranslateStore {
|
|||
return m[field]
|
||||
}
|
||||
|
||||
func (c *cluster) translateIndexKeys(ctx context.Context, index string, keys []string) ([]uint64, error) {
|
||||
keySet := make(map[string]struct{})
|
||||
for _, key := range keys {
|
||||
keySet[key] = struct{}{}
|
||||
}
|
||||
|
||||
keyMap, err := c.translateIndexKeySet(ctx, index, keySet)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
ids := make([]uint64, len(keys))
|
||||
for i := range keys {
|
||||
ids[i] = keyMap[keys[i]]
|
||||
}
|
||||
return ids, nil
|
||||
}
|
||||
|
||||
func (c *cluster) translateIndexKeySet(ctx context.Context, index string, keySet map[string]struct{}) (map[string]uint64, error) {
|
||||
keyMap := make(map[string]uint64)
|
||||
|
||||
// Split keys by partition.
|
||||
keysByPartition := make(map[int][]string, c.partitionN)
|
||||
for key := range keySet {
|
||||
partitionID := c.keyPartition(index, key)
|
||||
keysByPartition[partitionID] = append(keysByPartition[partitionID], key)
|
||||
}
|
||||
|
||||
// Translate keys by partition.
|
||||
var g errgroup.Group
|
||||
var mu sync.Mutex
|
||||
for partitionID := range keysByPartition {
|
||||
keys := keysByPartition[partitionID]
|
||||
|
||||
g.Go(func() (err error) {
|
||||
var ids []uint64
|
||||
if c.ownsPartition(c.Node.ID, partitionID) {
|
||||
if ids, err = c.translateIndexPartitionKeys(ctx, index, partitionID, keys); err != nil {
|
||||
return err
|
||||
}
|
||||
} else {
|
||||
nodes := c.partitionNodes(partitionID)
|
||||
if ids, err = c.InternalClient.TranslateKeysNode(ctx, &nodes[0].URI, index, "", keys); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
||||
mu.Lock()
|
||||
defer mu.Unlock()
|
||||
for i := range keys {
|
||||
keyMap[keys[i]] = ids[i]
|
||||
}
|
||||
return nil
|
||||
})
|
||||
}
|
||||
if err := g.Wait(); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return keyMap, nil
|
||||
}
|
||||
|
||||
func (c *cluster) translateIndexIDSet(ctx context.Context, index string, idSet map[uint64]struct{}) (map[uint64]string, error) {
|
||||
idMap := make(map[uint64]string)
|
||||
|
||||
// Split ids by partition.
|
||||
idsByPartition := make(map[int][]uint64, c.partitionN)
|
||||
for id := range idSet {
|
||||
partitionID := c.idPartition(index, id)
|
||||
idsByPartition[partitionID] = append(idsByPartition[partitionID], id)
|
||||
}
|
||||
|
||||
// Translate ids by partition.
|
||||
var g errgroup.Group
|
||||
var mu sync.Mutex
|
||||
for partitionID, ids := range idsByPartition {
|
||||
g.Go(func() error {
|
||||
nodes := c.partitionNodes(partitionID)
|
||||
keys, err := c.InternalClient.TranslateIDsNode(ctx, &nodes[0].URI, index, "", ids)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
mu.Lock()
|
||||
defer mu.Unlock()
|
||||
for i := range ids {
|
||||
idMap[ids[i]] = keys[i]
|
||||
}
|
||||
return nil
|
||||
})
|
||||
}
|
||||
if err := g.Wait(); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return idMap, nil
|
||||
}
|
||||
|
||||
func (c *cluster) translateIndexPartitionKeys(ctx context.Context, index string, partitionID int, keys []string) ([]uint64, error) {
|
||||
s := c.indexPartitionTranslateStore(index, partitionID)
|
||||
|
|
@ -2388,54 +2376,7 @@ func (c *cluster) setTranslateStoreReadOnly(v bool) {
|
|||
}
|
||||
}
|
||||
}
|
||||
|
||||
// translateEntryReader returns a reader that merges all index & field reader
|
||||
// that are specified in the offsets map.
|
||||
func (c *cluster) translateEntryReader(ctx context.Context, offsets TranslateOffsetMap) (_ TranslateEntryReader, err error) {
|
||||
// Ensure all readers are cleaned up if any error.
|
||||
var a []TranslateEntryReader
|
||||
defer func() {
|
||||
if err != nil {
|
||||
for i := range a {
|
||||
a[i].Close()
|
||||
}
|
||||
}
|
||||
}()
|
||||
|
||||
// Fetch all index partition readers.
|
||||
for indexName, indexMap := range offsets {
|
||||
for partitionID, offset := range indexMap.Partitions {
|
||||
store := c.indexPartitionTranslateStore(indexName, partitionID)
|
||||
if store == nil {
|
||||
return nil, ErrTranslateStoreNotFound
|
||||
}
|
||||
|
||||
r, err := store.EntryReader(ctx, uint64(offset))
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "index partition translate reader")
|
||||
}
|
||||
a = append(a, r)
|
||||
}
|
||||
}
|
||||
|
||||
// Fetch all field readers.
|
||||
for indexName, indexMap := range offsets {
|
||||
for fieldName, offset := range indexMap.Fields {
|
||||
store := c.fieldTranslateStore(indexName, fieldName)
|
||||
if store == nil {
|
||||
return nil, ErrTranslateStoreNotFound
|
||||
}
|
||||
|
||||
r, err := store.EntryReader(ctx, uint64(offset))
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "field translate reader")
|
||||
}
|
||||
a = append(a, r)
|
||||
}
|
||||
}
|
||||
|
||||
return NewMultiTranslateEntryReader(ctx, a), nil
|
||||
}
|
||||
*/
|
||||
|
||||
/*
|
||||
// holderTranslateStoreReplicator manages the replication of translation store
|
||||
|
|
|
|||
|
|
@ -273,6 +273,22 @@ func (Serializer) Unmarshal(buf []byte, m pilosa.Message) error {
|
|||
}
|
||||
decodeTranslateKeysResponse(msg, mt)
|
||||
return nil
|
||||
case *pilosa.TranslateIDsRequest:
|
||||
msg := &internal.TranslateIDsRequest{}
|
||||
err := proto.Unmarshal(buf, msg)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "unmarshaling TranslateIDsRequest")
|
||||
}
|
||||
decodeTranslateIDsRequest(msg, mt)
|
||||
return nil
|
||||
case *pilosa.TranslateIDsResponse:
|
||||
msg := &internal.TranslateIDsResponse{}
|
||||
err := proto.Unmarshal(buf, msg)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "unmarshaling TranslateIDsResponse")
|
||||
}
|
||||
decodeTranslateIDsResponse(msg, mt)
|
||||
return nil
|
||||
default:
|
||||
panic(fmt.Sprintf("unhandled pilosa.Message of type %T: %#v", mt, m))
|
||||
}
|
||||
|
|
@ -338,6 +354,10 @@ func encodeToProto(m pilosa.Message) proto.Message {
|
|||
return encodeTranslateKeysRequest(mt)
|
||||
case *pilosa.TranslateKeysResponse:
|
||||
return encodeTranslateKeysResponse(mt)
|
||||
case *pilosa.TranslateIDsRequest:
|
||||
return encodeTranslateIDsRequest(mt)
|
||||
case *pilosa.TranslateIDsResponse:
|
||||
return encodeTranslateIDsResponse(mt)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
|
|
|||
12
executor.go
12
executor.go
|
|
@ -3656,7 +3656,7 @@ func (e *executor) translateCall(index string, idx *Index, c *pql.Call, keyMap m
|
|||
return errors.Errorf("row value must be a string or non-negative integer, but got: %v of %[1]T", c.Args[rowKey])
|
||||
}
|
||||
} else if value := callArgString(c, rowKey); value != "" {
|
||||
id, err := e.Cluster.translateFieldKey(index, fieldName, value)
|
||||
id, err := field.TranslateStore().TranslateKey(value)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
|
@ -3752,7 +3752,7 @@ func (e *executor) translateGroupByCall(index string, idx *Index, c *pql.Call, k
|
|||
if !ok {
|
||||
return errors.New("prev value must be a string when field 'keys' option enabled")
|
||||
}
|
||||
id, err := e.Cluster.translateFieldKey(index, field.Name(), prevStr)
|
||||
id, err := field.TranslateStore().TranslateKey(prevStr)
|
||||
if err != nil {
|
||||
return errors.Wrapf(err, "translating row key '%s'", prevStr)
|
||||
}
|
||||
|
|
@ -3857,7 +3857,7 @@ func (e *executor) translateResult(index string, idx *Index, call *pql.Call, res
|
|||
return nil, fmt.Errorf("field %q not found", fieldName)
|
||||
}
|
||||
if field.keys() {
|
||||
key, err := field.translateStore.TranslateID(result.Pair.ID)
|
||||
key, err := field.TranslateStore().TranslateID(result.Pair.ID)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
|
@ -3881,7 +3881,7 @@ func (e *executor) translateResult(index string, idx *Index, call *pql.Call, res
|
|||
if field.keys() {
|
||||
other := make([]Pair, len(result.Pairs))
|
||||
for i := range result.Pairs {
|
||||
key, err := field.translateStore.TranslateID(result.Pairs[i].ID)
|
||||
key, err := field.TranslateStore().TranslateID(result.Pairs[i].ID)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
|
@ -3908,7 +3908,7 @@ func (e *executor) translateResult(index string, idx *Index, call *pql.Call, res
|
|||
return nil, ErrFieldNotFound
|
||||
}
|
||||
if field.keys() {
|
||||
key, err := e.Cluster.translateFieldID(index, g.Field, g.RowID)
|
||||
key, err := field.TranslateStore().TranslateID(g.RowID)
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "translating row ID in Group")
|
||||
}
|
||||
|
|
@ -3939,7 +3939,7 @@ func (e *executor) translateResult(index string, idx *Index, call *pql.Call, res
|
|||
} else if field.keys() {
|
||||
other.Keys = make([]string, len(result))
|
||||
for i, id := range result {
|
||||
key, err := e.Cluster.translateFieldID(index, fieldName, id)
|
||||
key, err := field.TranslateStore().TranslateID(id)
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "translating row ID")
|
||||
}
|
||||
|
|
|
|||
|
|
@ -54,10 +54,6 @@ func TestExecutor_TranslateGroupByCall(t *testing.T) {
|
|||
t.Fatalf("creating fields %v, %v, %v", erra, errb, errc)
|
||||
}
|
||||
|
||||
if err := cluster.updateTranslateStores(); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
query, err := pql.ParseString(`GroupBy(Rows(ak), Rows(b), Rows(ck), previous=["la", 0, "ha"], having=Condition(count > 10))`)
|
||||
if err != nil {
|
||||
t.Fatalf("parsing query: %v", err)
|
||||
|
|
|
|||
17
field.go
17
field.go
|
|
@ -89,6 +89,11 @@ type Field struct {
|
|||
logger logger.Logger
|
||||
|
||||
snapshotQueue snapshotQueue
|
||||
|
||||
translateStore TranslateStore
|
||||
|
||||
// Instantiates new translation stores
|
||||
OpenTranslateStore OpenTranslateStoreFunc
|
||||
}
|
||||
|
||||
// FieldOption is a functional option type for pilosa.fieldOptions.
|
||||
|
|
@ -260,6 +265,8 @@ func newField(path, index, name string, opts FieldOption) (*Field, error) {
|
|||
remoteAvailableShards: roaring.NewBitmap(),
|
||||
|
||||
logger: logger.NopLogger,
|
||||
|
||||
OpenTranslateStore: OpenInMemTranslateStore,
|
||||
}
|
||||
return f, nil
|
||||
}
|
||||
|
|
@ -278,6 +285,11 @@ func (f *Field) TranslateStorePath() string {
|
|||
return filepath.Join(f.path, "keys")
|
||||
}
|
||||
|
||||
// TranslateStore returns the field's translation store.
|
||||
func (f *Field) TranslateStore() TranslateStore {
|
||||
return f.translateStore
|
||||
}
|
||||
|
||||
// RowAttrStore returns the attribute storage.
|
||||
func (f *Field) RowAttrStore() AttrStore { return f.rowAttrStore }
|
||||
|
||||
|
|
@ -468,6 +480,11 @@ func (f *Field) Open() error {
|
|||
return errors.Wrap(err, "opening attrstore")
|
||||
}
|
||||
|
||||
f.logger.Debugf("open translate store for index/field: %s/%s", f.index, f.name)
|
||||
if f.translateStore, err = f.OpenTranslateStore(f.TranslateStorePath(), f.index, f.name, 0); err != nil {
|
||||
return errors.Wrap(err, "opening field translate store")
|
||||
}
|
||||
|
||||
return nil
|
||||
}(); err != nil {
|
||||
f.Close()
|
||||
|
|
|
|||
|
|
@ -79,6 +79,10 @@ type Holder struct {
|
|||
|
||||
// Manages replication from the primary node.
|
||||
primaryTranslateNode *Node
|
||||
|
||||
// Instantiates new translation stores
|
||||
OpenTranslateStore OpenTranslateStoreFunc
|
||||
OpenTranslateReader OpenTranslateReaderFunc
|
||||
}
|
||||
|
||||
// lockedChan looks a little ridiculous admittedly, but exists for good reason.
|
||||
|
|
|
|||
|
|
@ -1183,7 +1183,7 @@ func (c *InternalClient) TranslateIDsNode(ctx context.Context, uri *pilosa.URI,
|
|||
return nil, pilosa.ErrIndexRequired
|
||||
}
|
||||
|
||||
buf, err := c.serializer.Marshal(pilosa.TranslateIDsRequest{
|
||||
buf, err := c.serializer.Marshal(&pilosa.TranslateIDsRequest{
|
||||
Index: index,
|
||||
Field: field,
|
||||
IDs: ids,
|
||||
|
|
|
|||
|
|
@ -309,6 +309,7 @@ func newRouter(handler *Handler) *mux.Router {
|
|||
router.HandleFunc("/internal/index/{index}/attr/diff", handler.handlePostIndexAttrDiff).Methods("POST").Name("PostIndexAttrDiff")
|
||||
router.HandleFunc("/internal/translate/data", handler.handlePostTranslateData).Methods("POST").Name("PostTranslateData")
|
||||
router.HandleFunc("/internal/translate/keys", handler.handlePostTranslateKeys).Methods("POST").Name("PostTranslateKeys")
|
||||
router.HandleFunc("/internal/translate/ids", handler.handlePostTranslateIDs).Methods("POST").Name("PostTranslateIDs")
|
||||
router.HandleFunc("/internal/index/{index}/field/{field}/attr/diff", handler.handlePostFieldAttrDiff).Methods("POST").Name("PostFieldAttrDiff")
|
||||
router.HandleFunc("/internal/index/{index}/field/{field}/remote-available-shards/{shardID}", handler.handleDeleteRemoteAvailableShard).Methods("DELETE")
|
||||
router.HandleFunc("/internal/nodes", handler.handleGetNodes).Methods("GET").Name("GetNodes")
|
||||
|
|
@ -1757,6 +1758,30 @@ func (h *Handler) handlePostTranslateKeys(w http.ResponseWriter, r *http.Request
|
|||
http.Error(w, fmt.Sprintf("translate keys: %v", err), http.StatusInternalServerError)
|
||||
}
|
||||
|
||||
// Write response.
|
||||
_, err = w.Write(buf)
|
||||
if err != nil {
|
||||
h.logger.Printf("writing translate keys response: %v", err)
|
||||
return
|
||||
}
|
||||
}
|
||||
|
||||
func (h *Handler) handlePostTranslateIDs(w http.ResponseWriter, r *http.Request) {
|
||||
// Verify that request is only communicating over protobufs.
|
||||
if r.Header.Get("Content-Type") != "application/x-protobuf" {
|
||||
http.Error(w, "Unsupported media type", http.StatusUnsupportedMediaType)
|
||||
return
|
||||
} else if r.Header.Get("Accept") != "application/x-protobuf" {
|
||||
http.Error(w, "Not acceptable", http.StatusNotAcceptable)
|
||||
return
|
||||
}
|
||||
|
||||
buf, err := h.api.TranslateIDs(r.Context(), r.Body)
|
||||
if err != nil {
|
||||
http.Error(w, fmt.Sprintf("translate ids: %v", err), http.StatusInternalServerError)
|
||||
return
|
||||
}
|
||||
|
||||
// Write response.
|
||||
_, err = w.Write(buf)
|
||||
if err != nil {
|
||||
|
|
|
|||
26
index.go
26
index.go
|
|
@ -61,6 +61,12 @@ type Index struct {
|
|||
|
||||
// Used for notifying holder when a field is added.
|
||||
holder *Holder
|
||||
|
||||
// Per-partition translation stores
|
||||
translateStores map[int]TranslateStore
|
||||
|
||||
// Instantiates new translation stores
|
||||
OpenTranslateStore OpenTranslateStoreFunc
|
||||
}
|
||||
|
||||
// NewIndex returns a new instance of Index.
|
||||
|
|
@ -82,6 +88,10 @@ func NewIndex(path, name string) (*Index, error) {
|
|||
Stats: stats.NopStatsClient,
|
||||
logger: logger.NopLogger,
|
||||
trackExistence: true,
|
||||
|
||||
translateStores: make(map[int]TranslateStore),
|
||||
|
||||
OpenTranslateStore: OpenInMemTranslateStore,
|
||||
}, nil
|
||||
}
|
||||
|
||||
|
|
@ -96,6 +106,11 @@ func (i *Index) TranslateStorePath(partitionID int) string {
|
|||
return filepath.Join(i.path, "keys", strconv.Itoa(partitionID))
|
||||
}
|
||||
|
||||
// TranslateStore returns the translation store for a given partition.
|
||||
func (i *Index) TranslateStore(partitionID int) TranslateStore {
|
||||
return i.translateStores[partitionID]
|
||||
}
|
||||
|
||||
// Keys returns true if the index uses string keys.
|
||||
func (i *Index) Keys() bool { return i.keys }
|
||||
|
||||
|
|
@ -145,6 +160,16 @@ func (i *Index) Open() (err error) {
|
|||
return errors.Wrap(err, "opening attrstore")
|
||||
}
|
||||
|
||||
// TODO(BBJ): Support non-default partition counts.
|
||||
i.logger.Debugf("open translate store for index: %s", i.name)
|
||||
for partitionID := 0; partitionID < defaultPartitionN; partitionID++ {
|
||||
store, err := i.OpenTranslateStore(i.TranslateStorePath(partitionID), i.name, "", partitionID)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "opening index translate store")
|
||||
}
|
||||
i.translateStores[partitionID] = store
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
|
|
@ -445,6 +470,7 @@ func (i *Index) newField(path, name string) (*Field, error) {
|
|||
if i.snapshotQueue != nil {
|
||||
f.snapshotQueue = i.snapshotQueue
|
||||
}
|
||||
f.OpenTranslateStore = i.OpenTranslateStore
|
||||
return f, nil
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -284,7 +284,7 @@ func OptServerClusterHasher(h Hasher) ServerOption {
|
|||
// used to specify the translation data store type.
|
||||
func OptServerOpenTranslateStore(fn OpenTranslateStoreFunc) ServerOption {
|
||||
return func(s *Server) error {
|
||||
s.cluster.OpenTranslateStore = fn
|
||||
s.holder.OpenTranslateStore = fn
|
||||
return nil
|
||||
}
|
||||
}
|
||||
|
|
@ -293,7 +293,7 @@ func OptServerOpenTranslateStore(fn OpenTranslateStoreFunc) ServerOption {
|
|||
// used to specify the remote translation data reader.
|
||||
func OptServerOpenTranslateReader(fn OpenTranslateReaderFunc) ServerOption {
|
||||
return func(s *Server) error {
|
||||
s.cluster.OpenTranslateReader = fn
|
||||
s.holder.OpenTranslateReader = fn
|
||||
return nil
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -421,7 +421,7 @@ func (h grpcHandler) Inspect(req *pb.InspectRequest, stream pb.Pilosa_InspectSer
|
|||
|
||||
case "int":
|
||||
// Translate column key.
|
||||
id, err := index.TranslateStore().TranslateKey(col)
|
||||
id, err := h.api.TranslateIndexKey(context.Background(), index.Name(), col)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "translating column key")
|
||||
}
|
||||
|
|
@ -438,7 +438,7 @@ func (h grpcHandler) Inspect(req *pb.InspectRequest, stream pb.Pilosa_InspectSer
|
|||
}
|
||||
case "decimal":
|
||||
// Translate column key.
|
||||
id, err := index.TranslateStore().TranslateKey(col)
|
||||
id, err := h.api.TranslateIndexKey(context.Background(), index.Name(), col)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "translating column key")
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue