WIP: Thread OpenTranslateStore through Holder to Index

This commit is contained in:
Travis 2020-02-06 17:36:02 -06:00 • committed by Matt Jaffee
parent f9f6fce6b4
commit 49c8bf01a0
No known key found for this signature in database
GPG key ID: 08A3DFFF987B11BF
5 changed files with 37 additions and 6 deletions

View file

@ -144,7 +144,7 @@ func (s *TranslateStore) Size() int64 {
return tx.Size()
}
// TranslateKeys converts a string key to an integer ID.
// TranslateKey converts a string key to an integer ID.
// If key does not have an associated id then one is created.
func (s *TranslateStore) TranslateKey(key string) (id uint64, _ error) {
// Find id by key under read lock.
@ -188,8 +188,8 @@ func (s *TranslateStore) TranslateKey(key string) (id uint64, _ error) {
return id, nil
}
// TranslateKeys converts a string key to an integer ID.
// If key does not have an associated id then one is created.
// TranslateKeys converts a slice of string keys to a slice of integer IDs.
// If a key does not have an associated id then one is created.
func (s *TranslateStore) TranslateKeys(keys []string) (ids []uint64, _ error) {
if len(keys) == 0 {
return nil, nil

View file

@ -138,6 +138,8 @@ func NewHolder(partitionN int) *Holder {
cacheFlushInterval: defaultCacheFlushInterval,
OpenTranslateStore: OpenInMemTranslateStore,
Logger: logger.NopLogger,
}
}
@ -496,6 +498,7 @@ func (h *Holder) newIndex(path, name string) (*Index, error) {
index.newAttrStore = h.NewAttrStore
index.columnAttrs = h.NewAttrStore(filepath.Join(index.path, ".data"))
index.snapshotQueue = h.snapshotQueue
index.OpenTranslateStore = h.OpenTranslateStore
index.holder = h
return index, nil
}
@ -1011,7 +1014,7 @@ func (s *holderSyncer) initializeIndexTranslateReplication() error {
continue
}
// Connect to remote not and begin streaming.
// Connect to remote node and begin streaming.
rd, err := s.Holder.OpenTranslateReader(context.Background(), node.URI.String(), m)
if err != nil {
return err

View file

@ -107,7 +107,7 @@ func (i *Index) Path() string { return i.path }
// TranslateStorePath returns the translation database path for a partition.
func (i *Index) TranslateStorePath(partitionID int) string {
return filepath.Join(i.path, "keys", strconv.Itoa(partitionID))
return filepath.Join(i.path, translateStoreDir, strconv.Itoa(partitionID))
}
// TranslateStore returns the translation store for a given partition.

View file

@ -1031,13 +1031,30 @@ func TestHandler_Endpoints(t *testing.T) {
})
}
func TestCluster_TranslateStore(t *testing.T) {
cluster := make(test.Cluster, 1)
cluster[0] = test.NewCommandNode(true,
server.OptCommandServerOptions(
pilosa.OptServerOpenTranslateStore(boltdb.OpenTranslateStore),
pilosa.OptServerOpenTranslateReader(http.GetOpenTranslateReaderFunc(nil)),
),
)
cluster[0].Config.Gossip.Port = "0"
err := cluster[0].Start()
if err != nil {
t.Fatalf("starting cluster 0: %v", err)
}
test.MustDo("POST", cluster[0].URL()+"/index/i0", "{\"options\": {\"keys\": true}}")
}
func TestClusterTranslator(t *testing.T) {
cluster := make(test.Cluster, 2)
cluster[0] = test.NewCommandNode(true)
cluster[0].Config.Gossip.Port = "0"
err := cluster[0].Start()
if err != nil {
t.Fatalf("starting cluster 1: %v", err)
t.Fatalf("starting cluster 0: %v", err)
}
cluster[1] = test.NewCommandNode(false,
server.OptCommandServerOptions(

View file

@ -22,6 +22,12 @@ import (
"github.com/pkg/errors"
)
const (
// translateStoreDir is the subdirctory into which the partitioned
// translate store data is stored.
translateStoreDir = "_keys"
)
// Translate store errors.
var (
ErrTranslateStoreClosed = errors.New("translate store closed")
@ -70,6 +76,11 @@ type OpenTranslateStoreFunc func(path, index, field string, partitionID, partiti
// GenerateNextPartitionedID returns the next ID within the same partition.
func GenerateNextPartitionedID(index string, prev uint64, partitionID, partitionN int) uint64 {
// If the translation store is not partitioned, just return
// the next ID.
if partitionID == -1 {
return prev + 1
}
// Try to use the next ID if it is in the same partition.
// Otherwise find ID in next shard that has a matching partition.
for id := prev + 1; ; id += ShardWidth {