Merge pull request #60 from molecula/translation-sharding

Translation sharding
This commit is contained in:
Travis Turner 2020-01-31 10:44:51 -06:00 • committed by GitHub
commit b693688677
No known key found for this signature in database
GPG key ID: 4AEE18F83AFDEB23
38 changed files with 4324 additions and 4031 deletions

149
api.go
View file

@ -545,7 +545,8 @@ func (api *API) ExportCSV(ctx context.Context, indexName string, fieldName strin
var err error
if field.Keys() {
if rowStr, err = field.translateStore.TranslateID(rowID); err != nil {
// TODO: handle case: field.ForeignIndex
if rowStr, err = field.TranslateStore().TranslateID(rowID); err != nil {
return errors.Wrap(err, "translating row")
}
} else {
@ -553,7 +554,9 @@ func (api *API) ExportCSV(ctx context.Context, indexName string, fieldName strin
}
if index.Keys() {
if colStr, err = index.translateStore.TranslateID(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 +966,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 = field.translateStore.TranslateKeys(req.RowKeys); err != nil {
if req.RowIDs, err = field.TranslateStore().TranslateKeys(req.RowKeys); err != nil {
return errors.Wrap(err, "translating rows")
}
}
@ -973,7 +976,7 @@ func (api *API) Import(ctx context.Context, req *ImportRequest, opts ...ImportOp
if len(req.ColumnIDs) != 0 {
return errors.New("column ids cannot be used because index uses string keys")
}
if req.ColumnIDs, err = index.translateStore.TranslateKeys(req.ColumnKeys); err != nil {
if req.ColumnIDs, err = api.cluster.translateIndexKeys(ctx, req.Index, req.ColumnKeys); err != nil {
return errors.Wrap(err, "translating columns")
}
}
@ -1078,7 +1081,7 @@ func (api *API) ImportValue(ctx context.Context, req *ImportValueRequest, opts .
if len(req.ColumnIDs) != 0 {
return errors.New("column ids cannot be used because index uses string keys")
}
if req.ColumnIDs, err = index.translateStore.TranslateKeys(req.ColumnKeys); err != nil {
if req.ColumnIDs, err = api.cluster.translateIndexKeys(ctx, req.Index, req.ColumnKeys); err != nil {
return errors.Wrap(err, "translating columns")
}
req.Shard = math.MaxUint64
@ -1087,12 +1090,14 @@ func (api *API) ImportValue(ctx context.Context, req *ImportValueRequest, opts .
// Translate values when the field uses keys (for example, when
// the field has a ForeignIndex with keys).
if field.Keys() {
uints, err := field.translateStore.TranslateKeys(req.StringValues)
// Perform translation.
uints, err := api.cluster.translateIndexKeys(ctx, field.ForeignIndex(), req.StringValues)
if err != nil {
return errors.Wrap(err, "translating string values")
return err
}
// Because the BSI field supports negative value, we have to
// convert the slice of uint64 keys to a slice of int64.
// Because the BSI field supports negative values, we have to
// convert the uint64 keys to a slice of int64.
ints := make([]int64, len(uints))
for i := range uints {
ints[i] = int64(uints[i])
@ -1384,14 +1389,75 @@ 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.holder.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)
}
func (api *API) TranslateIndexIDs(ctx context.Context, indexName string, ids []uint64) ([]string, error) {
return api.cluster.translateIndexIDs(ctx, indexName, ids)
}
// TranslateKeys handles a TranslateKeyRequest.
func (api *API) TranslateKeys(r io.Reader) ([]byte, error) {
func (api *API) TranslateKeys(ctx context.Context, r io.Reader) (_ []byte, err error) {
var req TranslateKeysRequest
if buf, err := ioutil.ReadAll(r); err != nil {
return nil, NewBadRequestError(errors.Wrap(err, "read translate keys request error"))
@ -1400,13 +1466,22 @@ func (api *API) TranslateKeys(r io.Reader) ([]byte, error) {
}
// Lookup store for either index or field and translate keys.
store, err := api.holder.TranslateStore(req.Index, req.Field)
if err != nil {
return nil, err
}
ids, err := store.TranslateKeys(req.Keys)
if err != nil {
return nil, err
var ids []uint64
if req.Field == "" {
if ids, err = api.cluster.translateIndexKeys(ctx, req.Index, req.Keys); err != nil {
return nil, err
}
} else {
if field := api.holder.Field(req.Index, req.Field); field == nil {
return nil, ErrFieldNotFound
} else if fi := field.ForeignIndex(); fi != "" {
ids, err = api.cluster.translateIndexKeys(ctx, fi, req.Keys)
if err != nil {
return nil, err
}
} else if ids, err = field.TranslateStore().TranslateKeys(req.Keys); err != nil {
return nil, err
}
}
// Encode response.
@ -1417,6 +1492,42 @@ func (api *API) TranslateKeys(r io.Reader) ([]byte, error) {
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 fi := field.ForeignIndex(); fi != "" {
keys, err = api.cluster.translateIndexIDs(ctx, fi, req.IDs)
if err != nil {
return nil, err
}
} 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()

View file

@ -203,14 +203,15 @@ func TestAPI_Import(t *testing.T) {
// Generate some keyed records.
rowIDs := []uint64{}
colKeys := []string{}
timestamps := []int64{}
for i := 1; i <= 10; i++ {
rowIDs = append(rowIDs, rowID)
timestamps = append(timestamps, timestamp)
colKeys = append(colKeys, fmt.Sprintf("col%d", i))
}
// Keys are sharded so ordering is not guaranteed.
colKeys := []string{"col10", "col8", "col9", "col6", "col7", "col4", "col5", "col2", "col3", "col1"}
// Import data with keys to the coordinator (node0) and verify that it gets
// translated and forwarded to the owner of shard 0 (node1; because of offsetModHasher)
req := &pilosa.ImportRequest{
@ -261,63 +262,6 @@ func TestAPI_Import(t *testing.T) {
}
})
t.Run("RowKeyColumnID", func(t *testing.T) {
ctx := context.Background()
index := "rkci"
field := "f"
_, err := m0.API.CreateIndex(ctx, index, pilosa.IndexOptions{Keys: false})
if err != nil {
t.Fatalf("creating index: %v", err)
}
_, err = m0.API.CreateField(ctx, index, field, pilosa.OptFieldTypeSet(pilosa.DefaultCacheType, 100), pilosa.OptFieldKeys())
if err != nil {
t.Fatalf("creating field: %v", err)
}
rowKey := "rowkey"
// Generate some keyed records.
rowKeys := []string{rowKey, rowKey, rowKey}
colIDs := []uint64{1, 2, pilosa.ShardWidth + 1}
timestamps := []int64{0, 0, 0}
// Import data with keys to the coordinator (node0) and verify that it gets
// translated and forwarded to the owner of shard 0 (node1; because of offsetModHasher)
req := &pilosa.ImportRequest{
Index: index,
Field: field,
Shard: 0,
RowKeys: rowKeys,
ColumnIDs: colIDs,
Timestamps: timestamps,
}
if err := m0.API.Import(ctx, req); err != nil {
t.Fatal(err)
}
pql := fmt.Sprintf("Row(%s=%s)", field, rowKey)
// Query node0.
if res, err := m0.API.Query(ctx, &pilosa.QueryRequest{Index: index, Query: pql}); err != nil {
t.Fatal(err)
} else if columns := res.Results[0].(*pilosa.Row).Columns(); !reflect.DeepEqual(columns, colIDs) {
t.Fatalf("unexpected column ids: %+v", columns)
}
// Query node1.
if err := test.RetryUntil(5*time.Second, func() error {
if res, err := m1.API.Query(ctx, &pilosa.QueryRequest{Index: index, Query: pql}); err != nil {
return err
} else if columns := res.Results[0].(*pilosa.Row).Columns(); !reflect.DeepEqual(columns, colIDs) {
return fmt.Errorf("unexpected column ids: %+v", columns)
}
return nil
}); err != nil {
t.Fatal(err)
}
})
}
func TestAPI_ImportValue(t *testing.T) {
@ -356,12 +300,13 @@ func TestAPI_ImportValue(t *testing.T) {
// Generate some keyed records.
values := []int64{}
colKeys := []string{}
for i := 1; i <= 10; i++ {
values = append(values, int64(i))
colKeys = append(colKeys, fmt.Sprintf("col%d", i))
}
// Column keys are sharded so their order is not guaranteed.
colKeys := []string{"col10", "col8", "col9", "col6", "col7", "col4", "col5", "col2", "col3", "col1"}
// Import data with keys to the coordinator (node0) and verify that it gets
// translated and forwarded to the owner of shard 0 (node1; because of offsetModHasher)
req := &pilosa.ImportValueRequest{

View file

@ -32,8 +32,8 @@ var (
)
// OpenTranslateStore opens and initializes a boltdb translation store.
func OpenTranslateStore(path, index, field string) (pilosa.TranslateStore, error) {
s := NewTranslateStore(index, field)
func OpenTranslateStore(path, index, field string, partitionID, partitionN int) (pilosa.TranslateStore, error) {
s := NewTranslateStore(index, field, partitionID, partitionN)
s.Path = path
if err := s.Open(); err != nil {
return nil, err
@ -49,8 +49,10 @@ type TranslateStore struct {
mu sync.RWMutex
db *bolt.DB
index string
field string
index string
field string
partitionID int
partitionN int
once sync.Once
closing chan struct{}
@ -63,10 +65,12 @@ type TranslateStore struct {
}
// NewTranslateStore returns a new instance of TranslateStore.
func NewTranslateStore(index, field string) *TranslateStore {
func NewTranslateStore(index, field string, partitionID, partitionN int) *TranslateStore {
return &TranslateStore{
index: index,
field: field,
partitionID: partitionID,
partitionN: partitionN,
closing: make(chan struct{}),
writeNotify: make(chan struct{}),
}
@ -108,6 +112,11 @@ func (s *TranslateStore) Close() (err error) {
return nil
}
// PartitionID returns the partition id the store was initialized with.
func (s *TranslateStore) PartitionID() int {
return s.partitionID
}
// ReadOnly returns true if the store is in read-only mode.
func (s *TranslateStore) ReadOnly() bool {
s.mu.RLock()
@ -158,9 +167,10 @@ func (s *TranslateStore) TranslateKey(key string) (id uint64, _ error) {
bkt := tx.Bucket([]byte("keys"))
if id = findIDByKey(bkt, key); id != 0 {
return nil
} else if id, err = bkt.NextSequence(); err != nil {
return err
} else if err := bkt.Put([]byte(key), u64tob(id)); err != nil {
}
id = pilosa.GenerateNextPartitionedID(s.index, maxID(tx), s.partitionID, s.partitionN)
if err := bkt.Put([]byte(key), u64tob(id)); err != nil {
return err
} else if err := tx.Bucket([]byte("ids")).Put(u64tob(id), []byte(key)); err != nil {
return err
@ -220,9 +230,10 @@ func (s *TranslateStore) TranslateKeys(keys []string) (ids []uint64, _ error) {
if ids[i] = findIDByKey(bkt, key); ids[i] != 0 {
continue
} else if ids[i], err = bkt.NextSequence(); err != nil {
return err
} else if err := bkt.Put([]byte(key), u64tob(ids[i])); err != nil {
}
ids[i] = pilosa.GenerateNextPartitionedID(s.index, maxID(tx), s.partitionID, s.partitionN)
if err := bkt.Put([]byte(key), u64tob(ids[i])); err != nil {
return err
} else if err := tx.Bucket([]byte("ids")).Put(u64tob(ids[i]), []byte(key)); err != nil {
return err
@ -311,9 +322,7 @@ func (s *TranslateStore) notifyWrite() {
// MaxID returns the highest id in the store.
func (s *TranslateStore) MaxID() (max uint64, err error) {
if err := s.db.View(func(tx *bolt.Tx) error {
if key, _ := tx.Bucket([]byte("ids")).Cursor().Last(); key != nil {
max = btou64(key)
}
max = maxID(tx)
return nil
}); err != nil {
return 0, err
@ -321,6 +330,14 @@ func (s *TranslateStore) MaxID() (max uint64, err error) {
return max, nil
}
// MaxID returns the highest id in the store.
func maxID(tx *bolt.Tx) uint64 {
if key, _ := tx.Bucket([]byte("ids")).Cursor().Last(); key != nil {
return btou64(key)
}
return 0
}
type TranslateEntryReader struct {
ctx context.Context
store *TranslateStore

View file

@ -28,24 +28,24 @@ func TestTranslateStore_TranslateKey(t *testing.T) {
s := MustOpenNewTranslateStore()
defer MustCloseTranslateStore(s)
// Ensure initial key translates to ID 1.
// Ensure initial key translates to first ID for shard
if id, err := s.TranslateKey("foo"); err != nil {
t.Fatal(err)
} else if got, want := id, uint64(1); got != want {
} else if got, want := id, uint64(247463937); got != want {
t.Fatalf("TranslateKey()=%d, want %d", got, want)
}
// Ensure next key autoincrements.
if id, err := s.TranslateKey("bar"); err != nil {
t.Fatal(err)
} else if got, want := id, uint64(2); got != want {
} else if got, want := id, uint64(247463938); got != want {
t.Fatalf("TranslateKey()=%d, want %d", got, want)
}
// Ensure retranslating existing key returns original ID.
if id, err := s.TranslateKey("foo"); err != nil {
t.Fatal(err)
} else if got, want := id, uint64(1); got != want {
} else if got, want := id, uint64(247463937); got != want {
t.Fatalf("TranslateKey()=%d, want %d", got, want)
}
}
@ -57,29 +57,29 @@ func TestTranslateStore_TranslateKeys(t *testing.T) {
// Ensure initial keys translate to incrementing IDs.
if ids, err := s.TranslateKeys([]string{"foo", "bar"}); err != nil {
t.Fatal(err)
} else if got, want := ids[0], uint64(1); got != want {
} else if got, want := ids[0], uint64(247463937); got != want {
t.Fatalf("TranslateKeys()[0]=%d, want %d", got, want)
} else if got, want := ids[1], uint64(2); got != want {
} else if got, want := ids[1], uint64(247463938); got != want {
t.Fatalf("TranslateKeys()[1]=%d, want %d", got, want)
}
// Ensure retranslation returns original IDs.
if ids, err := s.TranslateKeys([]string{"foo", "bar"}); err != nil {
t.Fatal(err)
} else if got, want := ids[0], uint64(1); got != want {
} else if got, want := ids[0], uint64(247463937); got != want {
t.Fatalf("TranslateKeys()[0]=%d, want %d", got, want)
} else if got, want := ids[1], uint64(2); got != want {
} else if got, want := ids[1], uint64(247463938); got != want {
t.Fatalf("TranslateKeys()[1]=%d, want %d", got, want)
}
// Ensure retranslating with existing and non-existing keys returns correctly.
if ids, err := s.TranslateKeys([]string{"foo", "baz", "bar"}); err != nil {
t.Fatal(err)
} else if got, want := ids[0], uint64(1); got != want {
} else if got, want := ids[0], uint64(247463937); got != want {
t.Fatalf("TranslateKeys()[0]=%d, want %d", got, want)
} else if got, want := ids[1], uint64(3); got != want {
} else if got, want := ids[1], uint64(247463939); got != want {
t.Fatalf("TranslateKeys()[1]=%d, want %d", got, want)
} else if got, want := ids[2], uint64(2); got != want {
} else if got, want := ids[2], uint64(247463938); got != want {
t.Fatalf("TranslateKeys()[2]=%d, want %d", got, want)
}
}
@ -89,20 +89,23 @@ func TestTranslateStore_TranslateID(t *testing.T) {
defer MustCloseTranslateStore(s)
// Setup initial keys.
if _, err := s.TranslateKey("foo"); err != nil {
id1, err := s.TranslateKey("foo")
if err != nil {
t.Fatal(err)
} else if _, err := s.TranslateKey("bar"); err != nil {
}
id2, err := s.TranslateKey("bar")
if err != nil {
t.Fatal(err)
}
// Ensure IDs can be translated back to keys.
if key, err := s.TranslateID(1); err != nil {
if key, err := s.TranslateID(id1); err != nil {
t.Fatal(err)
} else if got, want := key, "foo"; got != want {
t.Fatalf("TranslateID()=%s, want %s", got, want)
}
if key, err := s.TranslateID(2); err != nil {
if key, err := s.TranslateID(id2); err != nil {
t.Fatal(err)
} else if got, want := key, "bar"; got != want {
t.Fatalf("TranslateID()=%s, want %s", got, want)
@ -114,12 +117,13 @@ func TestTranslateStore_TranslateIDs(t *testing.T) {
defer MustCloseTranslateStore(s)
// Setup initial keys.
if _, err := s.TranslateKeys([]string{"foo", "bar"}); err != nil {
ids, err := s.TranslateKeys([]string{"foo", "bar"})
if err != nil {
t.Fatal(err)
}
// Ensure IDs can be translated back to keys.
if keys, err := s.TranslateIDs([]uint64{1, 2, 3}); err != nil {
if keys, err := s.TranslateIDs([]uint64{ids[0], ids[1], 1}); err != nil {
t.Fatal(err)
} else if got, want := keys[0], "foo"; got != want {
t.Fatalf("TranslateIDs()[0]=%s, want %s", got, want)
@ -151,7 +155,7 @@ func TestTranslateStore_EntryReader(t *testing.T) {
// Read first entry.
if err := r.ReadEntry(&entry); err != nil {
t.Fatal(err)
} else if got, want := entry.ID, uint64(1); got != want {
} else if got, want := entry.ID, uint64(247463937); got != want {
t.Fatalf("ReadEntry() ID=%d, want %d", got, want)
} else if got, want := entry.Key, "foo"; got != want {
t.Fatalf("ReadEntry() Key=%s, want %s", got, want)
@ -160,7 +164,7 @@ func TestTranslateStore_EntryReader(t *testing.T) {
// Read next entry.
if err := r.ReadEntry(&entry); err != nil {
t.Fatal(err)
} else if got, want := entry.ID, uint64(2); got != want {
} else if got, want := entry.ID, uint64(247463938); got != want {
t.Fatalf("ReadEntry() ID=%d, want %d", got, want)
} else if got, want := entry.Key, "bar"; got != want {
t.Fatalf("ReadEntry() Key=%s, want %s", got, want)
@ -174,7 +178,7 @@ func TestTranslateStore_EntryReader(t *testing.T) {
// Read newly created entry.
if err := r.ReadEntry(&entry); err != nil {
t.Fatal(err)
} else if got, want := entry.ID, uint64(3); got != want {
} else if got, want := entry.ID, uint64(247463939); got != want {
t.Fatalf("ReadEntry() ID=%d, want %d", got, want)
} else if got, want := entry.Key, "baz"; got != want {
t.Fatalf("ReadEntry() Key=%s, want %s", got, want)
@ -211,7 +215,7 @@ func TestTranslateStore_EntryReader(t *testing.T) {
var entry pilosa.TranslateEntry
if err := r.ReadEntry(&entry); err != nil {
t.Fatal(err)
} else if got, want := entry.ID, uint64(1); got != want {
} else if got, want := entry.ID, uint64(247463937); got != want {
t.Fatalf("ReadEntry() ID=%d, want %d", got, want)
} else if got, want := entry.Key, "foo"; got != want {
t.Fatalf("ReadEntry() Key=%s, want %s", got, want)
@ -302,7 +306,7 @@ func MustNewTranslateStore() *boltdb.TranslateStore {
panic(err)
}
s := boltdb.NewTranslateStore("I", "F")
s := boltdb.NewTranslateStore("I", "F", 0, pilosa.DefaultPartitionN)
s.Path = f.Name()
return s
}

View file

@ -44,6 +44,8 @@ type FieldValue struct {
// While I understand that putting the entire Client behind an interface might require this many methods,
// I don't want to let it go unquestioned.
type InternalClient interface {
InternalQueryClient
MaxShardByIndex(ctx context.Context) (map[string]uint64, error)
Schema(ctx context.Context) ([]*IndexInfo, error)
PostSchema(ctx context.Context, uri *URI, s *Schema, remote bool) error
@ -51,7 +53,6 @@ type InternalClient interface {
FragmentNodes(ctx context.Context, index string, shard uint64) ([]*Node, error)
Nodes(ctx context.Context) ([]*Node, error)
Query(ctx context.Context, index string, queryRequest *QueryRequest) (*QueryResponse, error)
QueryNode(ctx context.Context, uri *URI, index string, queryRequest *QueryRequest) (*QueryResponse, error)
Import(ctx context.Context, index, field string, shard uint64, bits []Bit, opts ...ImportOption) error
ImportK(ctx context.Context, index, field string, bits []Bit, opts ...ImportOption) error
EnsureIndex(ctx context.Context, name string, options IndexOptions) error
@ -78,6 +79,8 @@ type InternalClient interface {
// InternalQueryClient is the internal interface for querying a node.
type InternalQueryClient interface {
QueryNode(ctx context.Context, uri *URI, index string, queryRequest *QueryRequest) (*QueryResponse, error)
TranslateKeysNode(ctx context.Context, uri *URI, index, field string, keys []string) ([]uint64, error)
TranslateIDsNode(ctx context.Context, uri *URI, index, field string, id []uint64) ([]string, error)
}
type nopInternalQueryClient struct{}
@ -86,6 +89,14 @@ func (n *nopInternalQueryClient) QueryNode(ctx context.Context, uri *URI, index
return nil, nil
}
func (n nopInternalQueryClient) TranslateKeysNode(ctx context.Context, uri *URI, index, field string, keys []string) ([]uint64, error) {
return nil, nil
}
func (n nopInternalQueryClient) TranslateIDsNode(ctx context.Context, uri *URI, index, field string, ids []uint64) ([]string, error) {
return nil, nil
}
func newNopInternalQueryClient() *nopInternalQueryClient {
return &nopInternalQueryClient{}
}
@ -125,6 +136,12 @@ func (n nopInternalClient) Query(ctx context.Context, index string, queryRequest
func (n nopInternalClient) QueryNode(ctx context.Context, uri *URI, index string, queryRequest *QueryRequest) (*QueryResponse, error) {
return nil, nil
}
func (n nopInternalClient) TranslateKeysNode(ctx context.Context, uri *URI, index, field string, keys []string) ([]uint64, error) {
return nil, nil
}
func (n nopInternalClient) TranslateIDsNode(ctx context.Context, uri *URI, index, field string, ids []uint64) ([]string, error) {
return nil, nil
}
func (n nopInternalClient) Import(ctx context.Context, index, field string, shard uint64, bits []Bit, opts ...ImportOption) error {
return nil
}

View file

@ -40,8 +40,8 @@ import (
)
const (
// defaultPartitionN is the default number of partitions in a cluster.
defaultPartitionN = 256
// DefaultPartitionN is the default number of partitions in a cluster.
DefaultPartitionN = 256
// ClusterState represents the state returned in the /status endpoint.
ClusterStateStarting = "STARTING"
@ -235,13 +235,15 @@ type cluster struct { // nolint: maligned
logger logger.Logger
InternalClient InternalClient
// OpenTranslateReader OpenTranslateReaderFunc
}
// newCluster returns a new instance of Cluster with defaults.
func newCluster() *cluster {
return &cluster{
Hasher: &jmphasher{},
partitionN: defaultPartitionN,
partitionN: DefaultPartitionN,
ReplicaN: 1,
joiningLeavingNodes: make(chan nodeAction, 10), // buffered channel
@ -390,14 +392,6 @@ func (c *cluster) addNode(node *Node) error {
return nil
}
// If the cluster membership has changed, reset the primary for
// translate store replication.
if c.holder != nil {
if err := c.holder.setPrimaryTranslateStore(c.unprotectedPrimaryReplicaNode()); err != nil {
return err
}
}
// add to topology
if c.Topology == nil {
return fmt.Errorf("Cluster.Topology is nil")
@ -417,14 +411,6 @@ func (c *cluster) removeNode(nodeID string) error {
// remove from cluster
c.removeNodeBasicSorted(nodeID)
// If the cluster membership has changed, reset the primary for
// translate store replication.
if c.holder != nil {
if err := c.holder.setPrimaryTranslateStore(c.unprotectedPrimaryReplicaNode()); err != nil {
return err
}
}
// remove from topology
if c.Topology == nil {
return fmt.Errorf("Cluster.Topology is nil")
@ -867,8 +853,12 @@ func (c *cluster) fragSources(to *cluster, idx *Index) (map[string][]*ResizeSour
return m, nil
}
// partition returns the partition that a shard belongs to.
func (c *cluster) partition(index string, shard uint64) int {
// shardPartition returns the partition that a shard belongs to.
func (c *cluster) shardPartition(index string, shard uint64) int {
return shardPartition(index, shard, c.partitionN)
}
func shardPartition(index string, shard uint64, partitionN int) int {
var buf [8]byte
binary.BigEndian.PutUint64(buf[:], shard)
@ -876,7 +866,25 @@ func (c *cluster) partition(index string, shard uint64) int {
h := fnv.New64a()
_, _ = h.Write([]byte(index))
_, _ = h.Write(buf[:])
return int(h.Sum64() % uint64(c.partitionN))
return int(h.Sum64() % uint64(partitionN))
}
// keyPartition returns the partition that a key belongs to.
func (c *cluster) keyPartition(index, key string) int {
return keyPartition(index, key, c.partitionN)
}
func keyPartition(index, key string, partitionN int) int {
// Hash the bytes and mod by partition count.
h := fnv.New64a()
_, _ = h.Write([]byte(index))
_, _ = h.Write([]byte(key))
return int(h.Sum64() % uint64(partitionN))
}
// idPartition returns the partition that an id belongs to.
func (c *cluster) idPartition(index string, id uint64) int {
return shardPartition(index, id/ShardWidth, c.partitionN)
}
// ShardNodes returns a list of nodes that own a fragment. Safe for concurrent use.
@ -888,7 +896,19 @@ func (c *cluster) ShardNodes(index string, shard uint64) []*Node {
// shardNodes returns a list of nodes that own a fragment. unprotected
func (c *cluster) shardNodes(index string, shard uint64) []*Node {
return c.partitionNodes(c.partition(index, shard))
return c.partitionNodes(c.shardPartition(index, shard))
}
// KeyNodes returns a list of nodes that own a fragment. Safe for concurrent use.
func (c *cluster) KeyNodes(index, key string) []*Node {
c.mu.RLock()
defer c.mu.RUnlock()
return c.keyNodes(index, key)
}
// keyNodes returns a list of nodes that own a key. unprotected
func (c *cluster) keyNodes(index, key string) []*Node {
return c.partitionNodes(c.keyPartition(index, key))
}
// ownsShard returns true if a host owns a fragment.
@ -900,7 +920,6 @@ func (c *cluster) ownsShard(nodeID string, index string, shard uint64) bool {
// partitionNodes returns a list of nodes that own a partition. unprotected.
func (c *cluster) partitionNodes(partitionID int) []*Node {
// Default replica count to between one and the number of nodes.
// The replica count can be zero if there are no nodes.
replicaN := c.ReplicaN
@ -922,11 +941,18 @@ func (c *cluster) partitionNodes(partitionID int) []*Node {
return nodes
}
// ownsPartition returns true if a host owns a partition.
func (c *cluster) ownsPartition(nodeID string, partition int) bool {
c.mu.RLock()
defer c.mu.RUnlock()
return Nodes(c.partitionNodes(partition)).ContainsID(nodeID)
}
// containsShards is like OwnsShards, but it includes replicas.
func (c *cluster) containsShards(index string, availableShards *roaring.Bitmap, node *Node) []uint64 {
var shards []uint64
availableShards.ForEach(func(i uint64) {
p := c.partition(index, i)
p := c.shardPartition(index, i)
// Determine the nodes for partition.
nodes := c.partitionNodes(p)
for _, n := range nodes {
@ -2047,6 +2073,148 @@ 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 {
partitionID := partitionID
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 := range idsByPartition {
partitionID := partitionID
ids := idsByPartition[partitionID]
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
}
// ClusterStatus describes the status of the cluster including its
// state and node topology.
type ClusterStatus struct {

View file

@ -96,7 +96,7 @@ func newIndexWithTempPath(name string) *Index {
if err != nil {
panic(err)
}
index, err := NewIndex(path, name)
index, err := NewIndex(path, name, DefaultPartitionN)
if err != nil {
panic(err)
}
@ -352,7 +352,7 @@ func TestCluster_Partition(t *testing.T) {
c := newCluster()
c.partitionN = partitionN
partitionID := c.partition(index, shard)
partitionID := c.shardPartition(index, shard)
if partitionID < 0 || partitionID >= partitionN {
t.Errorf("partition out of range: shard=%d, p=%d, n=%d", shard, partitionID, partitionN)
}

View file

@ -25,6 +25,7 @@ import (
"reflect"
"strings"
"testing"
"time"
"github.com/pilosa/pilosa/v2"
"github.com/pilosa/pilosa/v2/test"
@ -299,20 +300,26 @@ func TestImportCommand_KeyReplication(t *testing.T) {
// Verify that the data is available on both nodes.
for _, host := range []string{host0, host1} {
qry := "Count(Row(f=foo0))"
resp, err := http.DefaultClient.Do(MustNewHTTPRequest("POST", "http://"+host+"/index/i/query", strings.NewReader(qry)))
if err != nil {
t.Fatalf("Querying data for validation: %s", err)
}
if err := test.RetryUntil(2*time.Second, func() error {
qry := "Count(Row(f=foo0))"
resp, err := http.DefaultClient.Do(MustNewHTTPRequest("POST", "http://"+host+"/index/i/query", strings.NewReader(qry)))
if err != nil {
return fmt.Errorf("Querying data for validation: %s", err)
}
// Read body and unmarshal response.
exp := `{"results":[100]}` + "\n"
if body, err := ioutil.ReadAll(resp.Body); err != nil {
t.Fatalf("reading: %s", err)
} else if !reflect.DeepEqual(body, []byte(exp)) {
t.Fatalf("expected: %s, but got: %s", exp, body)
// Read body and unmarshal response.
exp := `{"results":[100]}` + "\n"
if body, err := ioutil.ReadAll(resp.Body); err != nil {
return fmt.Errorf("reading: %s", err)
} else if !reflect.DeepEqual(body, []byte(exp)) {
return fmt.Errorf("expected: %s, but got: %s", exp, body)
}
return nil
}); err != nil {
t.Fatal(err)
}
}
}
// Ensure that integer import with keys runs.

View file

@ -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
}
@ -760,17 +780,31 @@ func encodeRecalculateCaches(*pilosa.RecalculateCaches) *internal.RecalculateCac
return &internal.RecalculateCaches{}
}
func encodeTranslateKeysRequest(request *pilosa.TranslateKeysRequest) *internal.TranslateKeysRequest {
return &internal.TranslateKeysRequest{
Index: request.Index,
Field: request.Field,
Keys: request.Keys,
}
}
func encodeTranslateKeysResponse(response *pilosa.TranslateKeysResponse) *internal.TranslateKeysResponse {
return &internal.TranslateKeysResponse{
IDs: response.IDs,
}
}
func encodeTranslateKeysRequest(request *pilosa.TranslateKeysRequest) *internal.TranslateKeysRequest {
return &internal.TranslateKeysRequest{
func encodeTranslateIDsRequest(request *pilosa.TranslateIDsRequest) *internal.TranslateIDsRequest {
return &internal.TranslateIDsRequest{
Index: request.Index,
Field: request.Field,
Keys: request.Keys,
IDs: request.IDs,
}
}
func encodeTranslateIDsResponse(response *pilosa.TranslateIDsResponse) *internal.TranslateIDsResponse {
return &internal.TranslateIDsResponse{
Keys: response.Keys,
}
}
@ -1106,6 +1140,16 @@ func decodeTranslateKeysResponse(pb *internal.TranslateKeysResponse, m *pilosa.T
m.IDs = pb.IDs
}
func decodeTranslateIDsRequest(pb *internal.TranslateIDsRequest, m *pilosa.TranslateIDsRequest) {
m.Index = pb.Index
m.Field = pb.Field
m.IDs = pb.IDs
}
func decodeTranslateIDsResponse(pb *internal.TranslateIDsResponse, m *pilosa.TranslateIDsResponse) {
m.Keys = pb.Keys
}
// QueryResult types.
const (
queryResultTypeNil uint32 = iota

View file

@ -185,7 +185,7 @@ func (e *executor) Execute(ctx context.Context, index string, q *pql.Query, shar
// Translate query keys to ids, if necessary.
// No need to translate a remote call.
if !opt.Remote {
if err := e.translateCalls(ctx, index, idx, q.Calls); err != nil {
if err := e.translateCalls(ctx, index, q.Calls); err != nil {
return resp, err
} else if err := validateQueryContext(ctx); err != nil {
return resp, err
@ -230,12 +230,18 @@ func (e *executor) Execute(ctx context.Context, index string, q *pql.Query, shar
// Translate column attributes, if necessary.
if idx.Keys() {
idSet := make(map[uint64]struct{})
for _, col := range columnAttrSets {
v, err := idx.translateStore.TranslateID(col.ID)
if err != nil {
return resp, err
}
col.Key, col.ID = v, 0
idSet[col.ID] = struct{}{}
}
idMap, err := e.Cluster.translateIndexIDSet(ctx, index, idSet)
if err != nil {
return resp, errors.Wrap(err, "translating id set")
}
for _, col := range columnAttrSets {
col.Key, col.ID = idMap[col.ID], 0
}
}
@ -3520,69 +3526,118 @@ func (e *executor) mapperLocal(ctx context.Context, shards []uint64, mapFn mapFu
}
}
func (e *executor) translateCalls(ctx context.Context, index string, idx *Index, calls []*pql.Call) error {
func (e *executor) translateCalls(ctx context.Context, defaultIndexName string, calls []*pql.Call) (err error) {
span, _ := tracing.StartSpanFromContext(ctx, "Executor.translateCalls")
defer span.Finish()
// Generate a list of all used
keySets := make(map[string]map[string]struct{})
keySets[defaultIndexName] = make(map[string]struct{})
for i := range calls {
// Possibly change to another index for translation, if this
// call crosses index boundaries.
newIdxName := calls[i].CallIndex()
var newIdx *Index
if newIdxName == "" || newIdxName == index {
newIdxName = index
newIdx = idx
} else {
newIdx = idx.holder.indexes[newIdxName]
if newIdx == nil {
return fmt.Errorf("unknown index %q specified in cross-index call", newIdxName)
if err := e.collectCallKeySets(ctx, defaultIndexName, calls[i], keySets); err != nil {
return err
}
}
// Perform a separate batch translation for each separate index used.
keyMaps := make(map[string]map[string]uint64)
for indexName, keySet := range keySets {
idx := e.Holder.indexes[indexName]
if idx == nil {
return fmt.Errorf("cannot find index %q", indexName)
}
if !idx.Keys() || len(keySets) == 0 {
continue
}
if keyMaps[indexName], err = e.Cluster.translateIndexKeySet(ctx, indexName, keySet); err != nil {
return err
}
}
// Translate calls.
for i := range calls {
if err := e.translateCall(defaultIndexName, calls[i], keyMaps); err != nil {
return err
}
}
return nil
}
func (e *executor) collectCallKeySets(ctx context.Context, indexName string, c *pql.Call, m map[string]map[string]struct{}) error {
// Specifying an 'index' call overrides indexes on subsequent calls.
if s := c.CallIndex(); s != "" {
indexName = s
}
if m[indexName] == nil {
m[indexName] = make(map[string]struct{})
}
// Collect key for this call.
colKey, rowKey, fieldName := c.TranslateInfo(columnLabel, rowLabel)
if c.Args[colKey] != nil && isString(c.Args[colKey]) {
if value := callArgString(c, colKey); value != "" {
m[indexName][value] = struct{}{}
}
}
// Collect foreign index keys.
if fieldName != "" {
idx := e.Holder.indexes[indexName]
if field := idx.Field(fieldName); field != nil && field.ForeignIndex() != "" {
foreignIndexName := field.ForeignIndex()
if m[foreignIndexName] == nil {
m[foreignIndexName] = make(map[string]struct{})
}
if c.Args[rowKey] != nil && isCondition(c.Args[rowKey]) {
cond := c.Args[rowKey].(*pql.Condition)
if isString(cond.Value) {
m[foreignIndexName][cond.Value.(string)] = struct{}{}
}
} else if value := callArgString(c, rowKey); value != "" {
m[foreignIndexName][value] = struct{}{}
}
}
if err := e.translateCall(newIdxName, newIdx, calls[i]); err != nil {
}
// Recursively collect argument calls.
for _, arg := range c.Args {
if arg, ok := arg.(*pql.Call); ok {
if err := e.collectCallKeySets(ctx, indexName, arg, m); err != nil {
return errors.Wrap(err, "collecting group by call index name")
}
}
}
// Recursively collect child calls.
for _, child := range c.Children {
if err := e.collectCallKeySets(ctx, indexName, child, m); err != nil {
return err
}
}
return nil
}
func (e *executor) translateCall(index string, idx *Index, c *pql.Call) error {
var colKey, rowKey, fieldName string
switch c.Name {
case "Set", "Clear", "Row", "Range", "SetColumnAttrs", "ClearRow":
// Positional args in new PQL syntax require special handling here.
colKey = "_" + columnLabel
fieldName, _ = c.FieldArg()
rowKey = fieldName
case "SetRowAttrs":
// Positional args in new PQL syntax require special handling here.
rowKey = "_" + rowLabel
fieldName = callArgString(c, "_field")
case "Rows":
fieldName = callArgString(c, "_field")
rowKey = "previous"
colKey = "column"
case "GroupBy":
return errors.Wrap(e.translateGroupByCall(index, idx, c), "translating GroupBy")
case "IncludesColumn":
colKey = "column"
default:
colKey = "col"
fieldName = callArgString(c, "field")
rowKey = "row"
func (e *executor) translateCall(indexName string, c *pql.Call, keyMaps map[string]map[string]uint64) (err error) {
// Specifying an 'index' arg applies to all nested calls.
if s := c.CallIndex(); s != "" {
indexName = s
}
keyMap := keyMaps[indexName]
// Translate column key.
colKey, rowKey, fieldName := c.TranslateInfo(columnLabel, rowLabel)
idx := e.Holder.indexes[indexName]
if idx.Keys() {
if c.Args[colKey] != nil && !isString(c.Args[colKey]) {
if !isValidID(c.Args[colKey]) {
return errors.Errorf("column value must be a string or non-negative integer, but got: %v of %[1]T", c.Args[colKey])
}
} else if value := callArgString(c, colKey); value != "" {
id, err := idx.translateStore.TranslateKey(value)
if err != nil {
return err
}
c.Args[colKey] = id
c.Args[colKey] = keyMap[value]
}
} else {
if isString(c.Args[colKey]) {
@ -3620,8 +3675,47 @@ func (e *executor) translateCall(index string, idx *Index, c *pql.Call) error {
c.Args[rowKey] = rowID
}
} else if field.Keys() {
if err := e.translateRowKey(c, field.translateStore, rowKey); err != nil {
return errors.Wrap(err, "translating rowkey")
foreignIndexName := field.ForeignIndex()
if c.Args[rowKey] != nil && isCondition(c.Args[rowKey]) {
// In the case where a field has a foreign index with keys,
// allow `== "key"` or `!= "key"` to be used against the BSI
// field.
cond := c.Args[rowKey].(*pql.Condition)
if isString(cond.Value) {
switch cond.Op {
case pql.EQ, pql.NEQ:
var id uint64
if foreignIndexName != "" {
id = keyMaps[foreignIndexName][cond.Value.(string)]
} else {
if id, err = field.TranslateStore().TranslateKey(cond.Value.(string)); err != nil {
return errors.Wrap(err, "translating key")
}
}
c.Args[rowKey] = &pql.Condition{
Op: cond.Op,
Value: id,
}
default:
return errors.Errorf("conditional is not supported with string predicates: %s", cond.Op)
}
}
} else if c.Args[rowKey] != nil && !isString(c.Args[rowKey]) {
// allow passing row id directly (this can come in handy, but make sure it is a valid row id)
if !isValidID(c.Args[rowKey]) {
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 != "" {
var id uint64
if foreignIndexName != "" {
id = keyMaps[foreignIndexName][value]
} else {
if id, err = field.TranslateStore().TranslateKey(value); err != nil {
return err
}
}
c.Args[rowKey] = id
}
} else {
if isString(c.Args[rowKey]) {
@ -3632,135 +3726,65 @@ func (e *executor) translateCall(index string, idx *Index, c *pql.Call) error {
// Translate child calls.
for _, child := range c.Children {
// Possibly change to another index for translation, if this
// call crosses index boundaries.
newIdxName := child.CallIndex()
var newIdx *Index
if newIdxName == "" || newIdxName == index {
newIdxName = index
newIdx = idx
} else {
newIdx = idx.holder.indexes[newIdxName]
if newIdx == nil {
return fmt.Errorf("unknown index %q specified in cross-index call", newIdxName)
}
}
if err := e.translateCall(newIdxName, newIdx, child); err != nil {
if err := e.translateCall(indexName, child, keyMaps); err != nil {
return err
}
}
return nil
}
// Translate call args.
for _, arg := range c.Args {
if arg, ok := arg.(*pql.Call); ok {
if err := e.translateCall(indexName, arg, keyMaps); err != nil {
return errors.Wrap(err, "translating arg")
}
}
}
func (e *executor) translateRowKey(c *pql.Call, store TranslateStore, rowKey string) error {
if c.Args[rowKey] != nil && isCondition(c.Args[rowKey]) {
// In the case where a field has a foreign index with keys,
// allow `== "key"` or `!= "key"` to be used against the BSI
// field.
cond := c.Args[rowKey].(*pql.Condition)
if isString(cond.Value) {
switch cond.Op {
case pql.EQ, pql.NEQ:
id, err := store.TranslateKey(cond.Value.(string))
// GroupBy-specific call translation.
if c.Name == "GroupBy" {
prev, ok := c.Args["previous"]
if !ok {
return nil // nothing else to be translated
}
previous, ok := prev.([]interface{})
if !ok {
return errors.Errorf("'previous' argument must be list, but got %T", prev)
}
if len(c.Children) != len(previous) {
return errors.Errorf("mismatched lengths for previous: %d and children: %d in %s", len(previous), len(c.Children), c)
}
fields := make([]*Field, len(c.Children))
for i, child := range c.Children {
fieldname := callArgString(child, "_field")
field := idx.Field(fieldname)
if field == nil {
return errors.Wrapf(ErrFieldNotFound, "getting field '%s' from '%s'", fieldname, child)
}
fields[i] = field
}
for i, field := range fields {
prev := previous[i]
if field.Keys() {
prevStr, ok := prev.(string)
if !ok {
return errors.New("prev value must be a string when field 'keys' option enabled")
}
// TODO: does this need to take field.ForeignIndex() into consideration?
id, err := field.TranslateStore().TranslateKey(prevStr)
if err != nil {
return errors.Wrap(err, "translating key")
return errors.Wrapf(err, "translating row key '%s'", prevStr)
}
c.Args[rowKey] = &pql.Condition{
Op: cond.Op,
Value: id,
previous[i] = id
} else {
if prevStr, ok := prev.(string); ok {
return errors.Errorf("got string row val '%s' in 'previous' for field %s which doesn't use string keys", prevStr, field.Name())
}
default:
return errors.Errorf("conditional is not supported with string predicates: %s", cond.Op)
}
}
} else if c.Args[rowKey] != nil && !isString(c.Args[rowKey]) {
// allow passing row id directly (this can come in handy, but make sure it is a valid row id)
if !isValidID(c.Args[rowKey]) {
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 := store.TranslateKey(value)
if err != nil {
return errors.Wrap(err, "translating key")
}
c.Args[rowKey] = id
}
return nil
}
func (e *executor) translateGroupByCall(index string, idx *Index, c *pql.Call) error {
if c.Name != "GroupBy" {
panic("translateGroupByCall called with '" + c.Name + "'")
}
for _, child := range c.Children {
if err := e.translateCall(index, idx, child); err != nil {
return errors.Wrapf(err, "translating %s", child)
}
}
if filter, ok, err := c.CallArg("filter"); ok {
if err != nil {
return errors.Wrap(err, "getting filter call")
}
err = e.translateCall(index, idx, filter)
if err != nil {
return errors.Wrap(err, "translating filter call")
}
}
if aggregate, ok, err := c.CallArg("aggregate"); ok {
if err != nil {
return errors.Wrap(err, "getting aggregate call")
}
err = e.translateCall(index, idx, aggregate)
if err != nil {
return errors.Wrap(err, "translating aggregate call")
}
}
prev, ok := c.Args["previous"]
if !ok {
return nil // nothing else to be translated
}
previous, ok := prev.([]interface{})
if !ok {
return errors.Errorf("'previous' argument must be list, but got %T", prev)
}
if len(c.Children) != len(previous) {
return errors.Errorf("mismatched lengths for previous: %d and children: %d in %s", len(previous), len(c.Children), c)
}
fields := make([]*Field, len(c.Children))
for i, child := range c.Children {
fieldname := callArgString(child, "_field")
field := idx.Field(fieldname)
if field == nil {
return errors.Wrapf(ErrFieldNotFound, "getting field '%s' from '%s'", fieldname, child)
}
fields[i] = field
}
for i, field := range fields {
prev := previous[i]
if field.Keys() {
prevStr, ok := prev.(string)
if !ok {
return errors.New("prev value must be a string when field 'keys' option enabled")
}
id, err := field.translateStore.TranslateKey(prevStr)
if err != nil {
return errors.Wrapf(err, "translating row key '%s'", prevStr)
}
previous[i] = id
} else {
if prevStr, ok := prev.(string); ok {
return errors.Errorf("got string row val '%s' in 'previous' for field %s which doesn't use string keys", prevStr, field.Name())
}
}
}
return nil
}
@ -3768,8 +3792,22 @@ func (e *executor) translateResults(ctx context.Context, index string, idx *Inde
span, _ := tracing.StartSpanFromContext(ctx, "Executor.translateResults")
defer span.Finish()
idMap := make(map[uint64]string)
if idx.Keys() {
// Collect all index ids.
idSet := make(map[uint64]struct{})
for i := range calls {
if err := e.collectResultIDs(index, idx, calls[i], results[i], idSet); err != nil {
return err
}
}
if idMap, err = e.Cluster.translateIndexIDSet(ctx, index, idSet); err != nil {
return err
}
}
for i := range results {
results[i], err = e.translateResult(index, idx, calls[i], results[i])
results[i], err = e.translateResult(index, idx, calls[i], results[i], idMap)
if err != nil {
return err
}
@ -3777,18 +3815,30 @@ func (e *executor) translateResults(ctx context.Context, index string, idx *Inde
return nil
}
func (e *executor) translateResult(index string, idx *Index, call *pql.Call, result interface{}) (interface{}, error) {
func (e *executor) collectResultIDs(index string, idx *Index, call *pql.Call, result interface{}, idSet map[uint64]struct{}) error {
row, ok := result.(*Row)
if !ok {
return nil
} else if !idx.Keys() {
return nil
}
for _, segment := range row.Segments() {
for _, col := range segment.Columns() {
idSet[col] = struct{}{}
}
}
return nil
}
func (e *executor) translateResult(index string, idx *Index, call *pql.Call, result interface{}, idSet map[uint64]string) (interface{}, error) {
switch result := result.(type) {
case *Row:
if idx.Keys() {
other := &Row{Attrs: result.Attrs}
for _, segment := range result.Segments() {
for _, col := range segment.Columns() {
key, err := idx.translateStore.TranslateID(col)
if err != nil {
return nil, err
}
other.Keys = append(other.Keys, key)
other.Keys = append(other.Keys, idSet[col])
}
}
return other, nil
@ -3798,34 +3848,36 @@ func (e *executor) translateResult(index string, idx *Index, call *pql.Call, res
// make the return type for an int field with a ForeignIndex be
// a *Row instead (because it should always be positive).
case SignedRow:
var store TranslateStore
sr, err := func() (*SignedRow, error) {
fieldName := callArgString(call, "field")
if fieldName == "" {
return nil, nil
}
if fieldName := callArgString(call, "field"); fieldName != "" {
field := idx.Field(fieldName)
if field != nil && field.Keys() {
store = field.TranslateStore()
if field == nil {
return nil, nil
}
}
// In the case where a field/foreignIndex doesn't exist,
// fall back to using the index translateStore.
if store == nil && idx.Keys() {
store = idx.translateStore
}
if store != nil {
rslt := result.Pos
other := &Row{Attrs: rslt.Attrs}
for _, segment := range rslt.Segments() {
for _, col := range segment.Columns() {
key, err := store.TranslateID(col)
if field.Keys() {
rslt := result.Pos
other := &Row{Attrs: rslt.Attrs}
for _, segment := range rslt.Segments() {
keys, err := e.Cluster.translateIndexIDs(context.Background(), field.ForeignIndex(), segment.Columns())
if err != nil {
return nil, err
return nil, errors.Wrap(err, "translating index ids")
}
other.Keys = append(other.Keys, key)
other.Keys = append(other.Keys, keys...)
}
return &SignedRow{Pos: other}, nil
}
return SignedRow{Pos: other}, nil
return nil, nil
}()
if err != nil {
return nil, err
} else if sr != nil {
return *sr, nil
}
case PairField:
@ -3835,7 +3887,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
}
@ -3859,7 +3911,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
}
@ -3886,7 +3938,8 @@ func (e *executor) translateResult(index string, idx *Index, call *pql.Call, res
return nil, ErrFieldNotFound
}
if field.Keys() {
key, err := field.translateStore.TranslateID(g.RowID)
// TODO: does this need to take field.ForeignIndex() into consideration?
key, err := field.TranslateStore().TranslateID(g.RowID)
if err != nil {
return nil, errors.Wrap(err, "translating row ID in Group")
}
@ -3917,7 +3970,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 := field.translateStore.TranslateID(id)
key, err := field.TranslateStore().TranslateID(id)
if err != nil {
return nil, errors.Wrap(err, "translating row ID")
}

View file

@ -25,8 +25,13 @@ import (
)
func TestExecutor_TranslateGroupByCall(t *testing.T) {
holder := NewHolder(DefaultPartitionN)
cluster := NewTestCluster(1)
e := &executor{
Holder: NewHolder(),
Holder: holder,
Cluster: cluster,
}
e.Holder.Path, _ = ioutil.TempDir(*TempDir, "")
err := e.Holder.Open()
@ -51,7 +56,7 @@ func TestExecutor_TranslateGroupByCall(t *testing.T) {
t.Fatalf("parsing query: %v", err)
}
c := query.Calls[0]
err = e.translateGroupByCall("i", idx, c)
err = e.translateCall("i", c, make(map[string]map[string]uint64))
if err != nil {
t.Fatalf("translating call: %v", err)
}
@ -115,7 +120,7 @@ func TestExecutor_TranslateGroupByCall(t *testing.T) {
t.Fatalf("parsing query: %v", err)
}
c := query.Calls[0]
err = e.translateGroupByCall("i", idx, c)
err = e.translateCall("i", c, make(map[string]map[string]uint64))
if err == nil {
t.Fatalf("expected error, but translated call is '%s", c)
}

View file

@ -99,7 +99,7 @@ func TestExecutor_Execute_Row(t *testing.T) {
readQueries := []string{`Row(f=1)`}
responses := runCallTest(t, writeQuery, readQueries,
&pilosa.IndexOptions{Keys: true})
if keys := responses[0].Results[0].(*pilosa.Row).Keys; !reflect.DeepEqual(keys, []string{"one-hundred", "two-hundred"}) {
if keys := responses[0].Results[0].(*pilosa.Row).Keys; !reflect.DeepEqual(keys, []string{"two-hundred", "one-hundred"}) {
t.Fatalf("unexpected keys: %+v", keys)
}
})
@ -127,7 +127,7 @@ func TestExecutor_Execute_Row(t *testing.T) {
&pilosa.IndexOptions{Keys: true},
pilosa.OptFieldKeys())
if diff := cmp.Diff(responses[0].Results, []interface{}{
&pilosa.Row{Keys: []string{"foo", "bat"}, Attrs: map[string]interface{}{}},
&pilosa.Row{Keys: []string{"bat", "foo"}, Attrs: map[string]interface{}{}},
}, cmpopts.IgnoreUnexported(pilosa.Row{})); diff != "" {
t.Fatal(diff)
}
@ -324,7 +324,7 @@ func TestExecutor_Execute_Union(t *testing.T) {
readQueries := []string{`Union(Row(f=10), Row(f=11))`}
responses := runCallTest(t, writeQuery, readQueries,
&pilosa.IndexOptions{Keys: true})
if keys := responses[0].Results[0].(*pilosa.Row).Keys; !reflect.DeepEqual(keys, []string{"one", "one-hundred", "two-hundred", "two"}) {
if keys := responses[0].Results[0].(*pilosa.Row).Keys; !reflect.DeepEqual(keys, []string{"one", "two-hundred", "one-hundred", "two"}) {
t.Fatalf("unexpected keys: %+v", keys)
}
})
@ -357,7 +357,7 @@ func TestExecutor_Execute_Union(t *testing.T) {
responses := runCallTest(t, writeQuery, readQueries,
&pilosa.IndexOptions{Keys: true},
pilosa.OptFieldKeys())
if keys := responses[0].Results[0].(*pilosa.Row).Keys; !reflect.DeepEqual(keys, []string{"one", "one-hundred", "two-hundred", "two"}) {
if keys := responses[0].Results[0].(*pilosa.Row).Keys; !reflect.DeepEqual(keys, []string{"one", "two-hundred", "one-hundred", "two"}) {
t.Fatalf("unexpected keys: %+v", keys)
}
})
@ -1766,13 +1766,13 @@ func TestExecutor_Execute_Row_Range(t *testing.T) {
pilosa.OptFieldTypeTime(pilosa.TimeQuantum("YMDH")))
t.Run("Standard", func(t *testing.T) {
if keys := responses[0].Results[0].(*pilosa.Row).Keys; !reflect.DeepEqual(keys, []string{"two", "three", "four", "five", "six", "seven"}) {
if keys := responses[0].Results[0].(*pilosa.Row).Keys; !reflect.DeepEqual(keys, []string{"four", "six", "two", "seven", "five", "three"}) {
t.Fatalf("unexpected keys: %+v", keys)
}
})
t.Run("Clear", func(t *testing.T) {
if keys := responses[2].Results[0].(*pilosa.Row).Keys; !reflect.DeepEqual(keys, []string{"three", "four", "five", "six", "seven"}) {
if keys := responses[2].Results[0].(*pilosa.Row).Keys; !reflect.DeepEqual(keys, []string{"four", "six", "seven", "five", "three"}) {
t.Fatalf("unexpected keys: %+v", keys)
}
})
@ -1836,13 +1836,13 @@ func TestExecutor_Execute_Row_Range(t *testing.T) {
pilosa.OptFieldKeys())
t.Run("Standard", func(t *testing.T) {
if keys := responses[0].Results[0].(*pilosa.Row).Keys; !reflect.DeepEqual(keys, []string{"two", "three", "four", "five", "six", "seven"}) {
if keys := responses[0].Results[0].(*pilosa.Row).Keys; !reflect.DeepEqual(keys, []string{"four", "six", "two", "seven", "five", "three"}) {
t.Fatalf("unexpected keys: %+v", keys)
}
})
t.Run("Clear", func(t *testing.T) {
if keys := responses[2].Results[0].(*pilosa.Row).Keys; !reflect.DeepEqual(keys, []string{"three", "four", "five", "six", "seven"}) {
if keys := responses[2].Results[0].(*pilosa.Row).Keys; !reflect.DeepEqual(keys, []string{"four", "six", "seven", "five", "three"}) {
t.Fatalf("unexpected keys: %+v", keys)
}
})
@ -1950,13 +1950,13 @@ func TestExecutor_Execute_Range_Deprecated(t *testing.T) {
pilosa.OptFieldTypeTime(pilosa.TimeQuantum("YMDH")))
t.Run("Standard", func(t *testing.T) {
if keys := responses[0].Results[0].(*pilosa.Row).Keys; !reflect.DeepEqual(keys, []string{"two", "three", "four", "five", "six", "seven"}) {
if keys := responses[0].Results[0].(*pilosa.Row).Keys; !reflect.DeepEqual(keys, []string{"four", "six", "two", "seven", "five", "three"}) {
t.Fatalf("unexpected keys: %+v", keys)
}
})
t.Run("Clear", func(t *testing.T) {
if keys := responses[2].Results[0].(*pilosa.Row).Keys; !reflect.DeepEqual(keys, []string{"three", "four", "five", "six", "seven"}) {
if keys := responses[2].Results[0].(*pilosa.Row).Keys; !reflect.DeepEqual(keys, []string{"four", "six", "seven", "five", "three"}) {
t.Fatalf("unexpected keys: %+v", keys)
}
})
@ -2020,13 +2020,13 @@ func TestExecutor_Execute_Range_Deprecated(t *testing.T) {
pilosa.OptFieldKeys())
t.Run("Standard", func(t *testing.T) {
if keys := responses[0].Results[0].(*pilosa.Row).Keys; !reflect.DeepEqual(keys, []string{"two", "three", "four", "five", "six", "seven"}) {
if keys := responses[0].Results[0].(*pilosa.Row).Keys; !reflect.DeepEqual(keys, []string{"four", "six", "two", "seven", "five", "three"}) {
t.Fatalf("unexpected keys: %+v", keys)
}
})
t.Run("Clear", func(t *testing.T) {
if keys := responses[2].Results[0].(*pilosa.Row).Keys; !reflect.DeepEqual(keys, []string{"three", "four", "five", "six", "seven"}) {
if keys := responses[2].Results[0].(*pilosa.Row).Keys; !reflect.DeepEqual(keys, []string{"four", "six", "seven", "five", "three"}) {
t.Fatalf("unexpected keys: %+v", keys)
}
})
@ -2861,7 +2861,7 @@ func TestExecutor_Execute_Not(t *testing.T) {
TrackExistence: true,
Keys: true,
})
if keys := responses[0].Results[0].(*pilosa.Row).Keys; !reflect.DeepEqual(keys, []string{"three", "sw1"}) {
if keys := responses[0].Results[0].(*pilosa.Row).Keys; !reflect.DeepEqual(keys, []string{"sw1", "three"}) {
t.Fatalf("unexpected keys: %+v", keys)
}
})
@ -2892,7 +2892,7 @@ func TestExecutor_Execute_Not(t *testing.T) {
TrackExistence: true,
Keys: true,
}, pilosa.OptFieldKeys())
if keys := responses[0].Results[0].(*pilosa.Row).Keys; !reflect.DeepEqual(keys, []string{"three", "sw1"}) {
if keys := responses[0].Results[0].(*pilosa.Row).Keys; !reflect.DeepEqual(keys, []string{"sw1", "three"}) {
t.Fatalf("unexpected keys: %+v", keys)
}
})
@ -3016,12 +3016,12 @@ func TestExecutor_Execute_All(t *testing.T) {
expCols []string
expCnt uint64
}{
{qry: "All()", expCols: req.ColumnKeys, expCnt: uint64(bitCount)},
{qry: "All(limit=1)", expCols: req.ColumnKeys[:1], expCnt: 1},
{qry: "All(limit=4)", expCols: req.ColumnKeys, expCnt: 4},
{qry: "All(limit=5)", expCols: req.ColumnKeys, expCnt: 4},
{qry: "All(limit=1, offset=1)", expCols: req.ColumnKeys[1:2], expCnt: 1},
{qry: "All(limit=4, offset=1)", expCols: req.ColumnKeys[1:], expCnt: 3},
{qry: "All()", expCols: []string{"c1", "c0", "c3", "c2"}, expCnt: uint64(bitCount)},
{qry: "All(limit=1)", expCols: []string{"c1"}, expCnt: 1},
{qry: "All(limit=4)", expCols: []string{"c1", "c0", "c3", "c2"}, expCnt: 4},
{qry: "All(limit=5)", expCols: []string{"c1", "c0", "c3", "c2"}, expCnt: 4},
{qry: "All(limit=1, offset=1)", expCols: []string{"c0"}, expCnt: 1},
{qry: "All(limit=4, offset=1)", expCols: []string{"c0", "c3", "c2"}, expCnt: 3},
{qry: "All(limit=4, offset=5)", expCols: nil, expCnt: 0},
}
for i, test := range tests {
@ -3034,7 +3034,7 @@ func TestExecutor_Execute_All(t *testing.T) {
if len(cols) > 1000 || len(test.expCols) > 1000 {
t.Fatalf("test %d, unexpected columns, got: len(%d), but expected: len(%d)", i, len(cols), len(test.expCols))
} else {
t.Fatalf("test %d, unexpected columns, got: %T, but expected: %T", i, cols, test.expCols)
t.Fatalf("test %d, unexpected columns, got: %#v, but expected: %#v", i, cols, test.expCols)
}
}
}
@ -3712,7 +3712,7 @@ func TestExecutor_GroupByStrings(t *testing.T) {
Query: tst.query,
})
if err != nil {
t.Fatalf("got an error %v", err)
t.Fatal(err)
}
results := r.Results[0].([]pilosa.GroupCount)
test.CheckGroupBy(t, tst.expected, results)
@ -3898,7 +3898,7 @@ func TestExecutor_ForeignIndex(t *testing.T) {
`)
distinct := c.Query(t, "child", `Distinct(index="child", field="parent_id")`).Results[0].(pilosa.SignedRow)
if !reflect.DeepEqual(distinct.Pos.Keys, []string{"one", "two", "twenty-one"}) {
if !sameStringSlice(distinct.Pos.Keys, []string{"one", "two", "twenty-one"}) {
t.Fatalf("unexpected keys: %v", distinct.Pos.Keys)
}
@ -3918,6 +3918,31 @@ func TestExecutor_ForeignIndex(t *testing.T) {
}
}
// sameStringSlice is a helper function which compares two string
// slices without enforcing order.
func sameStringSlice(x, y []string) bool {
if len(x) != len(y) {
return false
}
// create a map of string -> int
diff := make(map[string]int, len(x))
for _, _x := range x {
// 0 value for int is 0, so just increment a counter for the string
diff[_x]++
}
for _, _y := range y {
// If the string _y is not in diff bail out early
if _, ok := diff[_y]; !ok {
return false
}
diff[_y] -= 1
if diff[_y] == 0 {
delete(diff, _y)
}
}
return len(diff) == 0
}
func TestExecutor_Execute_GroupBy(t *testing.T) {
groupByTest := func(t *testing.T, clusterSize int) {
c := test.MustRunCluster(t, 1)

103
field.go
View file

@ -75,9 +75,6 @@ type Field struct {
// Row attribute storage and cache
rowAttrStore AttrStore
// Key/ID translation store.
translateStore TranslateStore
broadcaster broadcaster
Stats stats.StatsClient
@ -101,7 +98,10 @@ type Field struct {
logger logger.Logger
snapshotQueue snapshotQueue
// Instantiates new translation store on open.
translateStore TranslateStore
// Instantiates new translation stores
OpenTranslateStore OpenTranslateStoreFunc
// Used for looking up a foreign index.
@ -314,12 +314,19 @@ func (f *Field) Index() string { return f.index }
// Path returns the path the field was initialized with.
func (f *Field) Path() string { return f.path }
// TranslateStorePath returns the translation database path for the field.
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 }
// TranslateStore returns the underlying translation store for the field.
func (f *Field) TranslateStore() TranslateStore { return f.translateStore }
// AvailableShards returns a bitmap of shards that contain data.
func (f *Field) AvailableShards() *roaring.Bitmap {
f.mu.RLock()
@ -513,16 +520,17 @@ func (f *Field) Open() error {
return errors.Wrap(err, "opening attrstore")
}
// If the field has a foreign index, and that index uses keys,
// then use that index's translateStore instead.
// Apply the field-specific translateStore.
if err := f.applyTranslateStore(); err != nil {
return errors.Wrap(err, "applying translate store")
}
// If the field has a foreign index, make sure the index
// exists.
if f.options.ForeignIndex != "" {
if err := f.holder.checkForeignIndex(f); err != nil {
return errors.Wrap(err, "checking foreign index")
}
} else {
if err := f.applyTranslateStore(); err != nil {
return errors.Wrap(err, "applying translate store")
}
}
return nil
@ -539,29 +547,43 @@ func (f *Field) Open() error {
func (f *Field) applyTranslateStore() error {
// Instantiate & open translation store.
var err error
f.translateStore, err = f.OpenTranslateStore(filepath.Join(f.path, "keys"), f.index, f.name)
f.translateStore, err = f.OpenTranslateStore(f.TranslateStorePath(), f.index, f.name, -1, -1)
if err != nil {
return errors.Wrap(err, "opening translate store")
return errors.Wrap(err, "opening field translate store")
}
f.usesKeys = f.options.Keys
// In the case where the field has a foreign index, set
// the usesKeys value accordingly.
if foreignIndexName := f.ForeignIndex(); foreignIndexName != "" {
if foreignIndex := f.holder.indexes[foreignIndexName]; foreignIndex != nil {
f.usesKeys = foreignIndex.Keys()
}
}
return nil
}
// applyForeignIndex sets the field's translateStore
// to that of a foreign index in the case where the
// foreign index uses keys. If the foreign index does
// not use keys, it falls back to applying the field's
// default translate store.
// applyForeignIndex used to set the field's translateStore to
// that of the foreign index, but since moving to partitioned
// translate stores on indexes, that doesn't happen anymore.
// So now all this method does is check that the foreign index
// actually exists. If we decided this was unnecessary (which
// it kind of is), we could remove the field.holder and all
// the logic which does this check on holder open after all
// indexes have opened.
func (f *Field) applyForeignIndex() error {
foreignIndex := f.holder.Index(f.options.ForeignIndex)
if foreignIndex == nil {
return errors.Wrapf(ErrForeignIndexNotFound, "%s", f.options.ForeignIndex)
} else if foreignIndex.Keys() {
f.usesKeys = true
f.translateStore = foreignIndex.translateStore
return nil
}
return f.applyTranslateStore()
f.usesKeys = foreignIndex.Keys()
return nil
}
// ForeignIndex returns the foreign index name attached to the field.
// Returns blank string if no foreign index exists.
func (f *Field) ForeignIndex() string {
return f.options.ForeignIndex
}
var fieldQueue = make(chan struct{}, 16)
@ -804,6 +826,13 @@ func (f *Field) Close() error {
_ = f.rowAttrStore.Close()
}
// Close field translation store.
if f.translateStore != nil {
if err := f.translateStore.Close(); err != nil {
return err
}
}
// Close all views.
for _, view := range f.viewMap {
if err := view.close(); err != nil {
@ -812,12 +841,6 @@ func (f *Field) Close() error {
}
f.viewMap = make(map[string]*view)
if f.translateStore != nil {
if err := f.translateStore.Close(); err != nil {
return err
}
}
return nil
}
@ -1583,18 +1606,20 @@ func (f *Field) importValue(columnIDs []uint64, values []int64, options *ImportO
requiredDepth = v
}
// Increase bit depth if required.
if requiredDepth > bsig.BitDepth {
if err := func() error {
f.mu.Lock()
defer f.mu.Unlock()
if err := func() error {
f.mu.Lock()
defer f.mu.Unlock()
bitDepth := bsig.BitDepth
if requiredDepth > bitDepth {
bsig.BitDepth = requiredDepth
f.options.BitDepth = requiredDepth
return f.saveMeta()
}(); err != nil {
return errors.Wrap(err, "increasing bsi bit depth")
} else {
requiredDepth = bitDepth
}
} else {
requiredDepth = bsig.BitDepth
return nil
}(); err != nil {
return errors.Wrap(err, "increasing bsi bit depth")
}
// Import into each fragment.

View file

@ -529,3 +529,50 @@ func TestField_ApplyOptions(t *testing.T) {
}
}
}
// Ensure that importValue handles requiredDepth correctly.
// This test sets the same column value to 1, then 8, then 1.
// A previous bug was incorrectly determining bitDepth based
// on the values in the import, and not taking existing values
// into consideration. This would cause an import of 1/8/1
// to result in a value of 9 instead of 1.
func TestBSIGroup_importValue(t *testing.T) {
f := MustOpenField(OptFieldTypeInt(-100, 200))
options := &ImportOptions{}
for i, tt := range []struct {
columnIDs []uint64
values []int64
checkVal int64
expCols []uint64
}{
{
[]uint64{100},
[]int64{1},
1,
[]uint64{100},
},
{
[]uint64{100},
[]int64{8},
8,
[]uint64{100},
},
{
[]uint64{100},
[]int64{1},
1,
[]uint64{100},
},
} {
if err := f.importValue(tt.columnIDs, tt.values, options); err != nil {
t.Fatalf("test %d, importing values: %s", i, err.Error())
}
if row, err := f.Range(f.name, pql.EQ, tt.checkVal); err != nil {
t.Fatalf("test %d, getting range: %s", i, err.Error())
} else if !reflect.DeepEqual(row.Columns(), tt.expCols) {
t.Fatalf("test %d, expected columns: %v, but got: %v", i, tt.expCols, row.Columns())
}
}
}

View file

@ -226,3 +226,17 @@ type TranslateKeysRequest struct {
type TranslateKeysResponse struct {
IDs []uint64
}
// TranslateIDsRequest describes the structure of a request
// for a batch of id translations.
type TranslateIDsRequest struct {
Index string
Field string
IDs []uint64
}
// TranslateIDsResponse is the structured response of a id
// translation request.
type TranslateIDsResponse struct {
Keys []string
}

470
holder.go
View file

@ -33,6 +33,7 @@ import (
"github.com/pilosa/pilosa/v2/tracing"
"github.com/pkg/errors"
uuid "github.com/satori/go.uuid"
"golang.org/x/sync/errgroup"
)
const (
@ -50,6 +51,9 @@ const (
type Holder struct {
mu sync.RWMutex
// Partition count used by translation.
partitionN int
// Indexes by name.
indexes map[string]*Index
@ -77,13 +81,9 @@ type Holder struct {
snapshotQueue snapshotQueue
// Manages replication from the primary node.
primaryTranslateNode *Node
translateStoreReplicator *holderTranslateStoreReplicator
// Instantiates new translation stores for indexes & fields.
OpenTranslateStore OpenTranslateStoreFunc // local store
OpenTranslateReader OpenTranslateReaderFunc // replication
// Instantiates new translation stores
OpenTranslateStore OpenTranslateStoreFunc
OpenTranslateReader OpenTranslateReaderFunc
// Queue of fields (having a foreign index) which have
// opened before their foreign index has opened.
@ -123,10 +123,11 @@ func (lc *lockedChan) Recv() {
}
// NewHolder returns a new instance of Holder.
func NewHolder() *Holder {
func NewHolder(partitionN int) *Holder {
return &Holder{
indexes: make(map[string]*Index),
closing: make(chan struct{}),
partitionN: partitionN,
indexes: make(map[string]*Index),
closing: make(chan struct{}),
opened: lockedChan{ch: make(chan struct{})},
@ -138,8 +139,6 @@ func NewHolder() *Holder {
cacheFlushInterval: defaultCacheFlushInterval,
Logger: logger.NopLogger,
OpenTranslateStore: OpenInMemTranslateStore,
}
}
@ -280,13 +279,6 @@ func (h *Holder) Close() error {
h.opened.ch = make(chan struct{})
h.opened.mu.Unlock()
h.mu.Lock()
if h.translateStoreReplicator != nil {
h.translateStoreReplicator.Close()
h.translateStoreReplicator = nil
}
h.mu.Unlock()
return nil
}
@ -500,14 +492,11 @@ func (h *Holder) createIndex(name string, opt IndexOptions) (*Index, error) {
// Update options.
h.indexes[index.Name()] = index
// Restart replication.
go h.refreshTranslateStoreReplicator()
return index, nil
}
func (h *Holder) newIndex(path, name string) (*Index, error) {
index, err := NewIndex(path, name)
index, err := NewIndex(path, name, h.partitionN)
if err != nil {
return nil, err
}
@ -518,7 +507,6 @@ func (h *Holder) newIndex(path, name string) (*Index, error) {
index.columnAttrs = h.NewAttrStore(filepath.Join(index.path, ".data"))
index.snapshotQueue = h.snapshotQueue
index.holder = h
index.OpenTranslateStore = h.OpenTranslateStore
return index, nil
}
@ -715,243 +703,6 @@ func (h *Holder) logStartup() error {
return nil
}
// TranslateStore returns store for the given index or field.
func (h *Holder) TranslateStore(index, field string) (TranslateStore, error) {
if field == "" {
idx := h.Index(index)
if idx == nil {
return nil, ErrIndexNotFound
}
return idx.TranslateStore(), nil
}
f := h.Field(index, field)
if f == nil {
return nil, ErrFieldNotFound
}
return f.TranslateStore(), nil
}
// TranslateOffsetMap returns a map of offsets for all indexes & fields.
func (h *Holder) TranslateOffsetMap() (TranslateOffsetMap, error) {
m := make(TranslateOffsetMap)
for _, idx := range h.Indexes() {
id, err := idx.TranslateStore().MaxID()
if err != nil {
return nil, err
}
m.SetIndexOffset(idx.Name(), id+1)
for _, field := range idx.Fields() {
id, err := field.TranslateStore().MaxID()
if err != nil {
return nil, err
}
m.SetFieldOffset(idx.Name(), field.Name(), id+1)
}
}
return m, nil
}
func (h *Holder) setTranslateStoreReadOnly(v bool) {
for _, idx := range h.Indexes() {
idx.TranslateStore().SetReadOnly(v)
for _, field := range idx.Fields() {
field.TranslateStore().SetReadOnly(v)
}
}
}
func (h *Holder) setPrimaryTranslateStore(node *Node) error {
if node != nil && h.OpenTranslateReader == nil {
return nil
}
h.mu.Lock()
h.primaryTranslateNode = node.Clone()
h.mu.Unlock()
go h.refreshTranslateStoreReplicator()
return nil
}
func (h *Holder) refreshTranslateStoreReplicator() {
h.mu.RLock()
node := h.primaryTranslateNode
h.mu.RUnlock()
var nodeURL string
if node != nil {
u := node.URI.URL()
nodeURL = u.String()
}
// Stop existing replication, if running.
h.mu.Lock()
if h.translateStoreReplicator != nil {
h.translateStoreReplicator.Close()
h.translateStoreReplicator = nil
}
h.mu.Unlock()
// Set all stores read only mode based on if we have a primary.
h.setTranslateStoreReadOnly(node != nil)
// Start replication monitor, if needed.
h.mu.Lock()
defer h.mu.Unlock()
if nodeURL != "" {
h.translateStoreReplicator = newHolderTranslateStoreReplicator(h, nodeURL)
h.translateStoreReplicator.logger = h.Logger
if err := h.translateStoreReplicator.Open(); err != nil {
h.Logger.Printf("cannot open translate store replicator: %s", err)
}
}
}
// TranslateEntryReader returns a reader that merges all index & field reader
// that are specified in the offsets map.
func (h *Holder) 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 readers.
for indexName, m := range offsets {
for fieldName, offset := range m {
var store TranslateStore
idx := h.Index(indexName)
if idx == nil {
return nil, ErrIndexNotFound
}
// Fetch from index or field store.
if fieldName == "" {
store = idx.TranslateStore()
} else {
f := idx.Field(fieldName)
if f == nil {
return nil, ErrFieldNotFound
}
store = f.TranslateStore()
}
// Generate reader and append to multireader.
r, err := store.EntryReader(ctx, uint64(offset))
if err != nil {
return nil, errors.Wrap(err, "translate reader")
}
a = append(a, r)
}
}
return NewMultiTranslateEntryReader(ctx, a), nil
}
// holderTranslateStoreReplicator manages the replication of translation store
// data from a primary store to the local replica. Continually tries to
// reconnect on disconnect.
type holderTranslateStoreReplicator struct {
ctx context.Context
cancel func()
wg sync.WaitGroup
holder *Holder
nodeURL string
logger logger.Logger
}
func newHolderTranslateStoreReplicator(h *Holder, nodeURL string) *holderTranslateStoreReplicator {
r := &holderTranslateStoreReplicator{
holder: h,
nodeURL: nodeURL,
logger: logger.NopLogger,
}
r.ctx, r.cancel = context.WithCancel(context.Background())
return r
}
// Open starts the background monitoring goroutine.
func (r *holderTranslateStoreReplicator) Open() error {
r.wg.Add(1)
go func() { defer r.wg.Done(); r.monitor() }()
return nil
}
// Close stops the replicator.
func (r *holderTranslateStoreReplicator) Close() error {
r.cancel()
return nil
}
// monitor runs in a background goroutine and continually tries to connect and
// stream translate changes from the primary store.
func (r *holderTranslateStoreReplicator) monitor() {
for {
select {
case <-r.ctx.Done():
return
default:
if err := r.replicate(); err != nil {
r.logger.Printf("cannot replicate: nodeURL=%s err=%s", r.nodeURL, err)
}
time.Sleep(1 * time.Second)
}
}
}
func (r *holderTranslateStoreReplicator) replicate() error {
// Determine the offsets of every index & field store.
offsets, err := r.holder.TranslateOffsetMap()
if err != nil {
return err
} else if len(offsets) == 0 {
return nil
}
// Begin streaming from remote primary.
rd, err := r.holder.OpenTranslateReader(r.ctx, r.nodeURL, offsets)
if err != nil {
return err
}
defer rd.Close()
for {
var entry TranslateEntry
if err := rd.ReadEntry(&entry); err != nil {
return err
}
// Find appropriate store.
var store TranslateStore
if entry.Field == "" {
idx := r.holder.Index(entry.Index)
if idx == nil {
return ErrIndexNotFound
}
store = idx.TranslateStore()
} else {
f := r.holder.Field(entry.Index, entry.Field)
if f == nil {
return ErrFieldNotFound
}
store = f.TranslateStore()
}
// Apply replication to store.
if err := store.ForceSet(entry.ID, entry.Key); err != nil {
return err
}
}
}
// holderSyncer is an active anti-entropy tool that compares the local holder
// with a remote holder based on block checksums and resolves differences.
type holderSyncer struct {
@ -962,6 +713,9 @@ type holderSyncer struct {
Node *Node
Cluster *cluster
// Translation sync handling.
readers []TranslateEntryReader
// Stats
Stats stats.StatsClient
@ -1175,6 +929,200 @@ func (s *holderSyncer) syncFragment(index, field, view string, shard uint64) err
return nil
}
// ResetTranslationSync reinitializes streaming sync of translation data.
func (s *holderSyncer) ResetTranslationSync() error {
// Stop existing streams.
if err := s.stopTranslationSync(); err != nil {
return errors.Wrap(err, "stop translation sync")
}
// Set read-only flag for all translation stores.
s.setTranslateReadOnlyFlags()
// Connect to each node that has a primary for which we are a replica.
if err := s.initializeIndexTranslateReplication(); err != nil {
return errors.Wrap(err, "initialize index translate replication")
}
// Connect to coordinator to stream field data.
if err := s.initializeFieldTranslateReplication(); err != nil {
return errors.Wrap(err, "initialize field translate replication")
}
return nil
}
// stopTranslationSync closes and waits for all outstanding translation readers
// to complete. This should be called before reconnecting to the cluster in case
// of a cluster resize or schema change.
func (s *holderSyncer) stopTranslationSync() error {
var g errgroup.Group
for i := range s.readers {
rd := s.readers[i]
g.Go(func() error {
return rd.Close()
})
}
return g.Wait()
}
// setTranslateReadOnlyFlags updates all translation stores to enabled or disable
// writing new translation keys. Index stores are writable if the node owns the
// partition. Field stores are writable if the node is the coordinator.
func (s *holderSyncer) setTranslateReadOnlyFlags() {
isCoordinator := s.Cluster.isCoordinator()
for _, index := range s.Holder.Indexes() {
for partitionID := 0; partitionID < s.Cluster.partitionN; partitionID++ {
ownsPartition := s.Cluster.ownsPartition(s.Node.ID, partitionID)
index.TranslateStore(partitionID).SetReadOnly(!ownsPartition)
}
for _, field := range index.Fields() {
field.TranslateStore().SetReadOnly(!isCoordinator)
}
}
}
// initializeIndexTranslateReplication connects to each node that is the
// primary for a partition that we are a replica of.
func (s *holderSyncer) initializeIndexTranslateReplication() error {
for _, node := range s.Cluster.Nodes() {
// Skip local node.
if node.ID == s.Node.ID {
continue
}
// Build a map of partition offsets to stream from.
m := make(TranslateOffsetMap)
for _, index := range s.Holder.Indexes() {
if !index.Keys() {
continue
}
for partitionID := 0; partitionID < s.Cluster.partitionN; partitionID++ {
partitionNodes := s.Cluster.partitionNodes(partitionID)
isPrimary := partitionNodes[0].ID == node.ID // remote is primary?
isReplica := Nodes(partitionNodes[1:]).ContainsID(s.Node.ID) // local is replica?
if !isPrimary || !isReplica {
continue
}
store := index.TranslateStore(partitionID)
offset, err := store.MaxID()
if err != nil {
return errors.Wrapf(err, "cannot determine max id for %q", index.Name())
}
m.SetIndexPartitionOffset(index.Name(), partitionID, offset)
}
}
// Skip if no replication required.
if len(m) == 0 {
continue
}
// Connect to remote not and begin streaming.
rd, err := s.Holder.OpenTranslateReader(context.Background(), node.URI.String(), m)
if err != nil {
return err
}
s.readers = append(s.readers, rd)
go func() { defer rd.Close(); s.readIndexTranslateReader(rd) }()
}
return nil
}
// initializeFieldTranslateReplication connects the coordinator to stream field data.
func (s *holderSyncer) initializeFieldTranslateReplication() error {
// Skip if coordinator.
if s.Cluster.Node.IsCoordinator {
return nil
}
// Build a map of partition offsets to stream from.
m := make(TranslateOffsetMap)
for _, index := range s.Holder.Indexes() {
if !index.Keys() {
continue
}
for _, field := range index.Fields() {
store := field.TranslateStore()
offset, err := store.MaxID()
if err != nil {
return errors.Wrapf(err, "cannot determine max id for %q/%q", index.Name(), field.Name())
}
m.SetFieldOffset(index.Name(), field.Name(), offset)
}
}
// Skip if no replication required.
if len(m) == 0 {
return nil
}
// Connect to coordinator and begin streaming.
coordinator := s.Cluster.coordinatorNode()
rd, err := s.Holder.OpenTranslateReader(context.Background(), coordinator.URI.String(), m)
if err != nil {
return err
}
s.readers = append(s.readers, rd)
go func() { defer rd.Close(); s.readFieldTranslateReader(rd) }()
return nil
}
func (s *holderSyncer) readIndexTranslateReader(rd TranslateEntryReader) {
for {
var entry TranslateEntry
if err := rd.ReadEntry(&entry); err != nil {
s.Holder.Logger.Printf("cannot read index translate entry: %s", err)
return
}
// Find appropriate store.
idx := s.Holder.Index(entry.Index)
if idx == nil {
s.Holder.Logger.Printf("index not found: %q", entry.Index)
return
}
// Apply replication to store.
store := idx.TranslateStore(s.Cluster.keyPartition(entry.Index, entry.Key))
if err := store.ForceSet(entry.ID, entry.Key); err != nil {
s.Holder.Logger.Printf("cannot force set index translation data: %d=%q", entry.ID, entry.Key)
return
}
}
}
func (s *holderSyncer) readFieldTranslateReader(rd TranslateEntryReader) {
for {
var entry TranslateEntry
if err := rd.ReadEntry(&entry); err != nil {
s.Holder.Logger.Printf("cannot read field translate entry: %s", err)
return
}
// Find appropriate store.
f := s.Holder.Field(entry.Index, entry.Field)
if f == nil {
s.Holder.Logger.Printf("field not found: %q/%q", entry.Index, entry.Field)
return
}
// Apply replication to store.
store := f.TranslateStore()
if err := store.ForceSet(entry.ID, entry.Key); err != nil {
s.Holder.Logger.Printf("cannot force set field translation data: %d=%q", entry.ID, entry.Key)
return
}
}
}
// holderCleaner removes fragments and data files that are no longer used.
type holderCleaner struct {
Node *Node

View file

@ -39,7 +39,7 @@ func (h *tHolder) Close() error {
// Note that the holder must be Closed first.
func (h *tHolder) Reopen() error {
path, logger := h.Path, h.Holder.Logger
h.Holder = NewHolder()
h.Holder = NewHolder(DefaultPartitionN)
h.Holder.Path = path
h.Holder.Logger = logger
return h.Holder.Open()
@ -51,7 +51,7 @@ func newHolder() *tHolder {
panic(err)
}
h := &tHolder{Holder: NewHolder()}
h := &tHolder{Holder: NewHolder(DefaultPartitionN)}
h.Path = path
return h
}
@ -302,7 +302,7 @@ func TestHolderCleaner_CleanHolder(t *testing.T) {
// Ensure holder can reopen.
func TestHolderCleaner_Reopen(t *testing.T) {
h := NewHolder()
h := NewHolder(DefaultPartitionN)
h.Path = "path"
err := h.Open()
if err != nil {

View file

@ -1124,6 +1124,106 @@ func (c *InternalClient) SendMessage(ctx context.Context, uri *pilosa.URI, msg [
return errors.Wrap(resp.Body.Close(), "closing response body")
}
// TranslateKeysNode sends a key translation request to a specific node.
func (c *InternalClient) TranslateKeysNode(ctx context.Context, uri *pilosa.URI, index, field string, keys []string) ([]uint64, error) {
span, ctx := tracing.StartSpanFromContext(ctx, "TranslateKeysNode")
defer span.Finish()
if index == "" {
return nil, pilosa.ErrIndexRequired
}
buf, err := c.serializer.Marshal(&pilosa.TranslateKeysRequest{
Index: index,
Field: field,
Keys: keys,
})
if err != nil {
return nil, errors.Wrap(err, "marshaling TranslateKeysRequest")
}
// Create HTTP request.
u := uri.Path("/internal/translate/keys")
req, err := http.NewRequest("POST", u, bytes.NewReader(buf))
if err != nil {
return nil, errors.Wrap(err, "creating request")
}
req.Header.Set("Content-Length", strconv.Itoa(len(buf)))
req.Header.Set("Content-Type", "application/x-protobuf")
req.Header.Set("Accept", "application/x-protobuf")
req.Header.Set("User-Agent", "pilosa/"+pilosa.Version)
// Execute request against the host.
resp, err := c.executeRequest(req.WithContext(ctx))
if err != nil {
return nil, err
}
defer resp.Body.Close()
// Read body and unmarshal response.
body, err := ioutil.ReadAll(resp.Body)
if err != nil {
return nil, errors.Wrap(err, "reading")
}
tkresp := &pilosa.TranslateKeysResponse{}
if err := c.serializer.Unmarshal(body, tkresp); err != nil {
return nil, fmt.Errorf("unmarshal response: %s", err)
}
return tkresp.IDs, nil
}
// TranslateIDsNode sends an id translation request to a specific node.
func (c *InternalClient) TranslateIDsNode(ctx context.Context, uri *pilosa.URI, index, field string, ids []uint64) ([]string, error) {
span, ctx := tracing.StartSpanFromContext(ctx, "TranslateIDsNode")
defer span.Finish()
if index == "" {
return nil, pilosa.ErrIndexRequired
}
buf, err := c.serializer.Marshal(&pilosa.TranslateIDsRequest{
Index: index,
Field: field,
IDs: ids,
})
if err != nil {
return nil, errors.Wrap(err, "marshaling TranslateIDsRequest")
}
// Create HTTP request.
u := uri.Path("/internal/translate/ids")
req, err := http.NewRequest("POST", u, bytes.NewReader(buf))
if err != nil {
return nil, errors.Wrap(err, "creating request")
}
req.Header.Set("Content-Length", strconv.Itoa(len(buf)))
req.Header.Set("Content-Type", "application/x-protobuf")
req.Header.Set("Accept", "application/x-protobuf")
req.Header.Set("User-Agent", "pilosa/"+pilosa.Version)
// Execute request against the host.
resp, err := c.executeRequest(req.WithContext(ctx))
if err != nil {
return nil, err
}
defer resp.Body.Close()
// Read body and unmarshal response.
body, err := ioutil.ReadAll(resp.Body)
if err != nil {
return nil, errors.Wrap(err, "reading")
}
tkresp := &pilosa.TranslateIDsResponse{}
if err := c.serializer.Unmarshal(body, tkresp); err != nil {
return nil, fmt.Errorf("unmarshal response: %s", err)
}
return tkresp.Keys, nil
}
// executeRequest executes the given request and checks the Response. For
// responses with non-2XX status, the body is read and closed, and an error is
// returned. If the error is nil, the caller must ensure that the response body

View file

@ -294,22 +294,29 @@ func TestClient_Export(t *testing.T) {
buf := bytes.NewBuffer(nil)
bw := bufio.NewWriter(buf)
// Send export request.
if err := c.ExportCSV(context.Background(), "keyed", "unkeyedf", 0, bw); err != nil {
t.Fatal(err)
// Send export request for every partition.
for i := 0; i < pilosa.DefaultPartitionN; i++ {
if err := c.ExportCSV(context.Background(), "keyed", "unkeyedf", uint64(i), bw); err != nil {
t.Fatal(err)
}
}
got := buf.String()
// Expected output.
exp := ""
for _, bit := range data {
exp += fmt.Sprintf("%d,%s\n", bit.RowID, bit.ColumnKey)
}
// Expected output is not sorted because of key sharding.
exp := "" +
"2,col200\n" +
"2,col201\n" +
"2,col202\n" +
"2,col203\n" +
"1,col103\n" +
"1,col102\n" +
"1,col101\n" +
"1,col100\n"
// Verify data.
if got != exp {
t.Fatalf("unexpected export data: %s", got)
t.Fatalf("unexpected export data: %q, expected %q", got, exp)
}
})
@ -329,21 +336,28 @@ func TestClient_Export(t *testing.T) {
bw := bufio.NewWriter(buf)
// Send export request.
if err := c.ExportCSV(context.Background(), "keyed", "keyedf", 0, bw); err != nil {
t.Fatal(err)
for i := 0; i < pilosa.DefaultPartitionN; i++ {
if err := c.ExportCSV(context.Background(), "keyed", "keyedf", uint64(i), bw); err != nil {
t.Fatal(err)
}
}
got := buf.String()
// Expected output.
exp := ""
for _, bit := range data {
exp += fmt.Sprintf("%s,%s\n", bit.RowKey, bit.ColumnKey)
}
// Expected output is unsorted because of key sharding.
exp := "" +
"row2,col200\n" +
"row2,col201\n" +
"row2,col202\n" +
"row2,col203\n" +
"row1,col103\n" +
"row1,col102\n" +
"row1,col101\n" +
"row1,col100\n"
// Verify data.
if got != exp {
t.Fatalf("unexpected export data: %s", got)
t.Fatalf("unexpected export data: %q, expected %q", got, exp)
}
})
}

View file

@ -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")
@ -1766,11 +1767,35 @@ func (h *Handler) handlePostTranslateKeys(w http.ResponseWriter, r *http.Request
return
}
buf, err := h.api.TranslateKeys(r.Body)
buf, err := h.api.TranslateKeys(r.Context(), r.Body)
if err != nil {
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 {

View file

@ -21,6 +21,7 @@ import (
"os"
"path/filepath"
"sort"
"strconv"
"sync"
"time"
@ -44,6 +45,9 @@ type Index struct {
trackExistence bool
existenceFld *Field
// Partitions used by translation.
partitionN int
// Fields by name.
fields map[string]*Field
@ -52,33 +56,34 @@ type Index struct {
// Column attribute storage and cache.
columnAttrs AttrStore
translateStore TranslateStore
broadcaster broadcaster
Stats stats.StatsClient
logger logger.Logger
snapshotQueue snapshotQueue
// Used for notifying holder when a field is added.
// Also passed to field for foreign-index lookup.
// Passed to field for foreign-index lookup.
holder *Holder
// Instantiates new translation stores for fields.
// Per-partition translation stores
translateStores map[int]TranslateStore
// Instantiates new translation stores
OpenTranslateStore OpenTranslateStoreFunc
}
// NewIndex returns a new instance of Index.
func NewIndex(path, name string) (*Index, error) {
func NewIndex(path, name string, partitionN int) (*Index, error) {
err := validateName(name)
if err != nil {
return nil, errors.Wrap(err, "validating name")
}
return &Index{
path: path,
name: name,
fields: make(map[string]*Field),
path: path,
name: name,
partitionN: partitionN,
fields: make(map[string]*Field),
newAttrStore: newNopAttrStore,
columnAttrs: nopStore,
@ -88,6 +93,8 @@ func NewIndex(path, name string) (*Index, error) {
logger: logger.NopLogger,
trackExistence: true,
translateStores: make(map[int]TranslateStore),
OpenTranslateStore: OpenInMemTranslateStore,
}, nil
}
@ -98,15 +105,22 @@ func (i *Index) Name() string { return i.name }
// Path returns the path the index was initialized with.
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))
}
// 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 }
// ColumnAttrStore returns the storage for column attributes.
func (i *Index) ColumnAttrStore() AttrStore { return i.columnAttrs }
// TranslateStore returns the underlying translation store for the index.
func (i *Index) TranslateStore() TranslateStore { return i.translateStore }
// Options returns all options for this index.
func (i *Index) Options() IndexOptions {
i.mu.RLock()
@ -150,9 +164,13 @@ func (i *Index) Open() (err error) {
return errors.Wrap(err, "opening attrstore")
}
// Instantiate & open translation store.
if i.translateStore, err = i.OpenTranslateStore(filepath.Join(i.path, "keys"), i.name, ""); err != nil {
return errors.Wrap(err, "opening translate store")
i.logger.Debugf("open translate store for index: %s", i.name)
for partitionID := 0; partitionID < i.partitionN; partitionID++ {
store, err := i.OpenTranslateStore(i.TranslateStorePath(partitionID), i.name, "", partitionID, i.partitionN)
if err != nil {
return errors.Wrap(err, "opening index translate store")
}
i.translateStores[partitionID] = store
}
return nil
@ -276,6 +294,14 @@ func (i *Index) Close() error {
// Close the attribute store.
i.columnAttrs.Close()
// Close partitioned translation stores.
for _, store := range i.translateStores {
if err := store.Close(); err != nil {
return errors.Wrap(err, "closing translation store")
}
}
i.translateStores = make(map[int]TranslateStore)
// Close all fields.
for _, f := range i.fields {
if err := f.Close(); err != nil {
@ -284,12 +310,6 @@ func (i *Index) Close() error {
}
i.fields = make(map[string]*Field)
if i.translateStore != nil {
if err := i.translateStore.Close(); err != nil {
return err
}
}
return nil
}
@ -450,11 +470,6 @@ func (i *Index) createField(name string, opt FieldOptions) (*Field, error) {
// Add to index's field lookup.
i.fields[name] = f
// Update replication, if needed.
if i.holder != nil {
go i.holder.refreshTranslateStoreReplicator()
}
return f, nil
}

View file

@ -25,7 +25,7 @@ func mustOpenIndex(opt IndexOptions) *Index {
if err != nil {
panic(err)
}
index, err := NewIndex(path, "i")
index, err := NewIndex(path, "i", DefaultPartitionN)
if err != nil {
panic(err)
}

View file

@ -217,7 +217,7 @@ func TestIndex_InvalidName(t *testing.T) {
if err != nil {
panic(err)
}
index, err := pilosa.NewIndex(path, "ABC")
index, err := pilosa.NewIndex(path, "ABC", pilosa.DefaultPartitionN)
if err == nil {
t.Fatalf("should have gotten an error on index name with caps")
}

File diff suppressed because it is too large Load diff

File diff suppressed because it is too large Load diff

View file

@ -133,6 +133,16 @@ message TranslateKeysResponse {
repeated uint64 IDs = 3;
}
message TranslateIDsRequest {
string Index = 1;
string Field = 2;
repeated uint64 IDs = 3;
}
message TranslateIDsResponse {
repeated string Keys = 3;
}
message ImportRoaringRequestView {
string Name = 1;
bytes Data = 2;

View file

@ -25,6 +25,7 @@ var _ pilosa.TranslateStore = (*TranslateStore)(nil)
type TranslateStore struct {
CloseFunc func() error
MaxIDFunc func() (uint64, error)
PartitionIDFunc func() int
ReadOnlyFunc func() bool
SetReadOnlyFunc func(v bool)
TranslateKeyFunc func(key string) (uint64, error)
@ -43,6 +44,10 @@ func (s *TranslateStore) MaxID() (uint64, error) {
return s.MaxIDFunc()
}
func (s *TranslateStore) PartitionID() int {
return s.PartitionIDFunc()
}
func (s *TranslateStore) ReadOnly() bool {
return s.ReadOnlyFunc()
}

View file

@ -338,6 +338,8 @@ var callInfoByFunc = map[string]callInfo{
"Row": {allowUnknown: true},
"Range": {allowUnknown: true},
"Distinct": {allowUnknown: true},
// allow only "field=X" cases with string field names
"Max": allowField,
"Min": allowField,
@ -742,6 +744,36 @@ func (c *Call) HasConditionArg() bool {
return false
}
// TranslateInfo returns the relevant translation fields.
func (c *Call) TranslateInfo(columnLabel, rowLabel string) (colKey, rowKey, fieldName string) {
switch c.Name {
case "Set", "Clear", "Row", "Range", "SetColumnAttrs", "ClearRow":
// Positional args in new PQL syntax require special handling here.
fieldName, _ = c.FieldArg()
return "_" + columnLabel, fieldName, fieldName
case "SetRowAttrs":
// Positional args in new PQL syntax require special handling here.
return "", "_" + rowLabel, c.ArgString("_field")
case "Rows":
return "column", "previous", c.ArgString("_field")
case "IncludesColumn":
return "column", "", ""
case "GroupBy":
return "", "", ""
default:
return "col", "row", c.ArgString("_field")
}
}
func (c *Call) ArgString(key string) string {
value, ok := c.Args[key]
if !ok {
return ""
}
s, _ := value.(string)
return s
}
// Condition represents an operation & value.
// When used in an argument map it represents a binary expression.
type Condition struct {

View file

@ -300,10 +300,12 @@ func OptServerOpenTranslateReader(fn OpenTranslateReaderFunc) ServerOption {
// NewServer returns a new instance of Server.
func NewServer(opts ...ServerOption) (*Server, error) {
cluster := newCluster()
s := &Server{
closing: make(chan struct{}),
cluster: newCluster(),
holder: NewHolder(),
cluster: cluster,
holder: NewHolder(cluster.partitionN),
diagnostics: newDiagnosticsCollector(defaultDiagnosticServer),
systemInfo: newNopSystemInfo(),
defaultClient: nopInternalClient{},
@ -527,6 +529,10 @@ func (s *Server) Open() error {
s.syncer.Closing = s.closing
s.syncer.Stats = s.holder.Stats.WithTags("HolderSyncer")
if err := s.syncer.ResetTranslationSync(); err != nil {
return err
}
// Start background monitoring.
s.wg.Add(3)
go func() { defer s.wg.Done(); s.monitorAntiEntropy() }()
@ -660,10 +666,16 @@ func (s *Server) receiveMessage(m Message) error {
if err != nil {
return err
}
if err := s.syncer.ResetTranslationSync(); err != nil {
return err
}
case *DeleteIndexMessage:
if err := s.holder.DeleteIndex(obj.Index); err != nil {
return err
}
if err := s.syncer.ResetTranslationSync(); err != nil {
return err
}
case *CreateFieldMessage:
idx := s.holder.Index(obj.Index)
if idx == nil {
@ -674,11 +686,17 @@ func (s *Server) receiveMessage(m Message) error {
if err != nil {
return err
}
if err := s.syncer.ResetTranslationSync(); err != nil {
return err
}
case *DeleteFieldMessage:
idx := s.holder.Index(obj.Index)
if err := idx.DeleteField(obj.Field); err != nil {
return err
}
if err := s.syncer.ResetTranslationSync(); err != nil {
return err
}
case *DeleteAvailableShardMessage:
f := s.holder.Field(obj.Index, obj.Field)
if err := f.RemoveAvailableShard(obj.ShardID); err != nil {

View file

@ -256,10 +256,31 @@ func (h grpcHandler) Inspect(req *pb.InspectRequest, stream pb.Pilosa_InspectSer
case "int":
if field.Keys() {
value, exists, err := field.StringValue(col)
if err != nil {
return errors.Wrap(err, "getting string field value for column")
} else if exists {
var value string
var exists bool
var err error
if fi := field.ForeignIndex(); fi != "" {
// Get the value from the int field.
intVal, ok, err := field.Value(col)
if err != nil {
return errors.Wrap(err, "getting int value")
} else if ok {
vals, err := h.api.TranslateIndexIDs(context.Background(), fi, []uint64{uint64(intVal)})
if err != nil {
return errors.Wrap(err, "getting keys for ids")
}
if len(vals) > 0 && vals[0] != "" {
value = vals[0]
exists = true
}
}
} else {
value, exists, err = field.StringValue(col)
if err != nil {
return errors.Wrap(err, "getting string field value for column")
}
}
if exists {
rowResp.Columns = append(rowResp.Columns,
&pb.ColumnResponse{ColumnVal: &pb.ColumnResponse_StringVal{StringVal: value}})
} else {
@ -454,16 +475,37 @@ 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")
}
if field.Keys() {
value, exists, err := field.StringValue(id)
if err != nil {
return errors.Wrap(err, "getting string field value for column")
} else if exists {
var value string
var exists bool
var err error
if fi := field.ForeignIndex(); fi != "" {
// Get the value from the int field.
intVal, ok, err := field.Value(id)
if err != nil {
return errors.Wrap(err, "getting int value")
} else if ok {
vals, err := h.api.TranslateIndexIDs(context.Background(), fi, []uint64{uint64(intVal)})
if err != nil {
return errors.Wrap(err, "getting keys for ids")
}
if len(vals) > 0 && vals[0] != "" {
value = vals[0]
exists = true
}
}
} else {
value, exists, err = field.StringValue(id)
if err != nil {
return errors.Wrap(err, "getting string field value for column")
}
}
if exists {
rowResp.Columns = append(rowResp.Columns,
&pb.ColumnResponse{ColumnVal: &pb.ColumnResponse_StringVal{StringVal: value}})
} else {
@ -485,7 +527,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")
}

View file

@ -991,7 +991,7 @@ func TestHandler_Endpoints(t *testing.T) {
if w.Code != gohttp.StatusOK {
t.Fatalf("unexpected status code: %d", w.Code)
}
target := []uint64{1, 2, 3}
target := []uint64{162529281, 159383553, 160432129}
resp := pilosa.TranslateKeysResponse{}
err = cmd.API.Serializer.Unmarshal(w.Body.Bytes(), &resp)
if err != nil {

View file

@ -137,7 +137,6 @@ func NewCommand(stdin io.Reader, stdout, stderr io.Writer, opts ...CommandOption
// Start starts the pilosa server - it returns once the server is running.
func (m *Command) Start() (err error) {
// Seed random number generator
rand.Seed(time.Now().UTC().UnixNano())

View file

@ -35,7 +35,7 @@ func NewHolder() *Holder {
panic(err)
}
h := &Holder{Holder: pilosa.NewHolder()}
h := &Holder{Holder: pilosa.NewHolder(pilosa.DefaultPartitionN)}
h.Path = path
h.Holder.NewAttrStore = boltdb.NewAttrStore
@ -61,7 +61,7 @@ func (h *Holder) Close() error {
// Note that the holder must be Closed first.
func (h *Holder) Reopen() error {
path, logger := h.Path, h.Holder.Logger
h.Holder = pilosa.NewHolder()
h.Holder = pilosa.NewHolder(pilosa.DefaultPartitionN)
h.Holder.Path = path
h.Holder.Logger = logger
h.Holder.NewAttrStore = boltdb.NewAttrStore

View file

@ -32,7 +32,7 @@ func newIndex() *Index {
if err != nil {
panic(err)
}
index, err := pilosa.NewIndex(path, "i")
index, err := pilosa.NewIndex(path, "i", pilosa.DefaultPartitionN)
if err != nil {
panic(err)
}
@ -62,7 +62,7 @@ func (i *Index) Reopen() error {
}
path, name := i.Path(), i.Name()
i.Index, err = pilosa.NewIndex(path, name)
i.Index, err = pilosa.NewIndex(path, name, pilosa.DefaultPartitionN)
if err != nil {
return err
}

View file

@ -137,6 +137,7 @@ func (m *Command) Reopen() error {
// MustCreateIndex uses this command's API to create an index and fails the test
// if there is an error.
func (m *Command) MustCreateIndex(tb testing.TB, name string, opts pilosa.IndexOptions) *pilosa.Index {
tb.Helper()
idx, err := m.API.CreateIndex(context.Background(), name, opts)
if err != nil {
tb.Fatalf("creating index: %v with options: %v, err: %v", name, opts, err)
@ -147,6 +148,7 @@ func (m *Command) MustCreateIndex(tb testing.TB, name string, opts pilosa.IndexO
// MustCreateField uses this command's API to create the field. The index must
// already exist - it fails the test if there is an error.
func (m *Command) MustCreateField(tb testing.TB, index, field string, opts ...pilosa.FieldOption) *pilosa.Field {
tb.Helper()
f, err := m.API.CreateField(context.Background(), index, field, opts...)
if err != nil {
tb.Fatalf("creating field: %s in index: %s err: %v", field, index, err)
@ -157,6 +159,7 @@ func (m *Command) MustCreateField(tb testing.TB, index, field string, opts ...pi
// MustQuery uses this command's API to execute the given query request, failing
// if Query returns a non-nil error, otherwise returning the QueryResponse.
func (m *Command) MustQuery(tb testing.TB, req *pilosa.QueryRequest) pilosa.QueryResponse {
tb.Helper()
resp, err := m.API.Query(context.Background(), req)
if err != nil {
tb.Fatalf("making query: %v, err: %v", req, err)
@ -248,6 +251,7 @@ type Cluster []*Command
// Query executes an API.Query through one of the cluster's node's API. It fails
// the test if there is an error.
func (c Cluster) Query(t testing.TB, index, query string) pilosa.QueryResponse {
t.Helper()
if len(c) == 0 {
t.Fatal("must have at least one node in cluster to query")
}
@ -256,6 +260,7 @@ func (c Cluster) Query(t testing.TB, index, query string) pilosa.QueryResponse {
}
func (c Cluster) ImportBits(t testing.TB, index, field string, rowcols [][2]uint64) {
t.Helper()
byShard := make(map[uint64][][2]uint64)
for _, rowcol := range rowcols {
shard := rowcol[1] / pilosa.ShardWidth
@ -296,6 +301,7 @@ func (c Cluster) ImportBits(t testing.TB, index, field string, rowcols [][2]uint
// CreateField creates the index (if necessary) and field specified.
func (c Cluster) CreateField(t testing.TB, index string, iopts pilosa.IndexOptions, field string, fopts ...pilosa.FieldOption) *pilosa.Field {
t.Helper()
idx, err := c[0].API.CreateIndex(context.Background(), index, iopts)
if err != nil && !strings.Contains(err.Error(), "index already exists") {
t.Fatalf("creating index: %v", err)
@ -344,6 +350,7 @@ func (c Cluster) Close() error {
// MustNewCluster creates a new cluster
func MustNewCluster(tb testing.TB, size int, opts ...[]server.CommandOption) Cluster {
tb.Helper()
c, err := newCluster(size, opts...)
if err != nil {
tb.Fatalf("new cluster: %v", err)
@ -391,6 +398,7 @@ func runCluster(size int, opts ...[]server.CommandOption) (Cluster, error) {
// MustRunCluster creates and starts a new cluster
func MustRunCluster(tb testing.TB, size int, opts ...[]server.CommandOption) Cluster {
tb.Helper()
c, err := runCluster(size, opts...)
if err != nil {
tb.Fatalf("run cluster: %v", err)
@ -429,6 +437,7 @@ func MustDo(method, urlStr string, body string) *httpResponse {
}
func CheckGroupBy(t *testing.T, expected, results []pilosa.GroupCount) {
t.Helper()
if len(results) != len(expected) {
t.Fatalf("number of groupings mismatch:\n got:%+v\nwant:%+v\n", results, expected)
}

View file

@ -28,6 +28,7 @@ var (
ErrTranslateStoreReaderClosed = errors.New("translate store reader closed")
ErrReplicationNotSupported = errors.New("replication not supported")
ErrTranslateStoreReadOnly = errors.New("translate store could not find or create key, translate store read only")
ErrTranslateStoreNotFound = errors.New("translate store not found")
ErrCannotOpenV1TranslateFile = errors.New("cannot open v1 translate .keys file")
)
@ -38,11 +39,18 @@ type TranslateStore interface {
// Returns the maximum ID set on the store.
MaxID() (uint64, error)
// Retrieves the partition ID associated with the store.
// Only applies to index stores.
PartitionID() int
// Sets & retrieves whether the store is read-only.
ReadOnly() bool
SetReadOnly(v bool)
// Converts a string key to its autoincrementing integer ID value.
//
// Translated id must be associated with a shard in the store's partition
// unless partition is set to -1.
TranslateKey(key string) (uint64, error)
TranslateKeys(key []string) ([]uint64, error)
@ -58,7 +66,18 @@ type TranslateStore interface {
}
// OpenTranslateStoreFunc represents a function for instantiating and opening a TranslateStore.
type OpenTranslateStoreFunc func(path, index, field string) (TranslateStore, error)
type OpenTranslateStoreFunc func(path, index, field string, partitionID, partitionN int) (TranslateStore, error)
// GenerateNextPartitionedID returns the next ID within the same partition.
func GenerateNextPartitionedID(index string, prev uint64, partitionID, partitionN int) uint64 {
// 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 {
if shardPartition(index, id/ShardWidth, partitionN) == partitionID {
return id
}
}
}
// TranslateEntryReader represents a stream of translation entries.
type TranslateEntryReader interface {
@ -154,22 +173,22 @@ type readEntryResponse struct {
}
// TranslateOffsetMap maintains a set of offsets for both indexes & fields.
type TranslateOffsetMap map[string]map[string]uint64
type TranslateOffsetMap map[string]*IndexTranslateOffsetMap
// IndexOffset returns the offset for the given index.
func (m TranslateOffsetMap) IndexOffset(name string) uint64 {
func (m TranslateOffsetMap) IndexPartitionOffset(name string, partitionID int) uint64 {
if m[name] == nil {
return 0
}
return m[name][""]
return m[name].Partitions[partitionID]
}
// SetIndexOffset sets the offset for the given index.
func (m TranslateOffsetMap) SetIndexOffset(name string, offset uint64) {
func (m TranslateOffsetMap) SetIndexPartitionOffset(name string, partitionID int, offset uint64) {
if m[name] == nil {
m[name] = make(map[string]uint64)
m[name] = NewIndexTranslateOffsetMap()
}
m[name][""] = offset
m[name].Partitions[partitionID] = offset
}
// FieldOffset returns the offset for the given field.
@ -177,15 +196,27 @@ func (m TranslateOffsetMap) FieldOffset(index, name string) uint64 {
if m[index] == nil {
return 0
}
return m[index][name]
return m[index].Fields[name]
}
// SetFieldOffset sets the offset for the given field.
func (m TranslateOffsetMap) SetFieldOffset(index, name string, offset uint64) {
if m[index] == nil {
m[index] = make(map[string]uint64)
m[index] = NewIndexTranslateOffsetMap()
}
m[index].Fields[name] = offset
}
type IndexTranslateOffsetMap struct {
Partitions map[int]uint64 `json:"partitions"`
Fields map[string]uint64 `json:"fields"`
}
func NewIndexTranslateOffsetMap() *IndexTranslateOffsetMap {
return &IndexTranslateOffsetMap{
Partitions: make(map[int]uint64),
Fields: make(map[string]uint64),
}
m[index][name] = offset
}
// Ensure type implements interface.
@ -193,22 +224,28 @@ var _ TranslateStore = &InMemTranslateStore{}
// InMemTranslateStore is an in-memory storage engine for mapping keys to int values.
type InMemTranslateStore struct {
mu sync.RWMutex
index string
field string
readOnly bool
keys []string
lookup map[string]uint64
mu sync.RWMutex
index string
field string
partitionID int
partitionN int
readOnly bool
keysByID map[uint64]string
idsByKey map[string]uint64
maxID uint64
writeNotify chan struct{}
}
// NewInMemTranslateStore returns a new instance of InMemTranslateStore.
func NewInMemTranslateStore(index, field string) *InMemTranslateStore {
func NewInMemTranslateStore(index, field string, partitionID, partitionN int) *InMemTranslateStore {
return &InMemTranslateStore{
index: index,
field: field,
lookup: make(map[string]uint64),
partitionID: partitionID,
partitionN: partitionN,
keysByID: make(map[uint64]string),
idsByKey: make(map[string]uint64),
writeNotify: make(chan struct{}),
}
}
@ -217,14 +254,19 @@ var _ OpenTranslateStoreFunc = OpenInMemTranslateStore
// OpenInMemTranslateStore returns a new instance of InMemTranslateStore.
// Implements OpenTranslateStoreFunc.
func OpenInMemTranslateStore(rawurl, index, field string) (TranslateStore, error) {
return NewInMemTranslateStore(index, field), nil
func OpenInMemTranslateStore(rawurl, index, field string, partitionID, partitionN int) (TranslateStore, error) {
return NewInMemTranslateStore(index, field, partitionID, partitionN), nil
}
func (s *InMemTranslateStore) Close() error {
return nil
}
// PartitionID returns the partition id the store was initialized with.
func (s *InMemTranslateStore) PartitionID() int {
return s.partitionID
}
// ReadOnly returns true if the store is in read-only mode.
func (s *InMemTranslateStore) ReadOnly() bool {
s.mu.Lock()
@ -244,40 +286,41 @@ func (s *InMemTranslateStore) SetReadOnly(v bool) {
func (s *InMemTranslateStore) TranslateKey(key string) (uint64, error) {
s.mu.Lock()
defer s.mu.Unlock()
if s.readOnly {
return 0, nil
}
return s.translateKey(key), nil
return s.translateKey(key)
}
// TranslateKeys converts a string key to an integer ID.
// If key does not have an associated id then one is created.
func (s *InMemTranslateStore) TranslateKeys(keys []string) ([]uint64, error) {
func (s *InMemTranslateStore) TranslateKeys(keys []string) (_ []uint64, err error) {
s.mu.Lock()
defer s.mu.Unlock()
if s.readOnly {
return nil, nil
}
ids := make([]uint64, len(keys))
for i := range keys {
ids[i] = s.translateKey(keys[i])
if ids[i], err = s.translateKey(keys[i]); err != nil {
return ids, err
}
}
return ids, nil
}
func (s *InMemTranslateStore) translateKey(key string) uint64 {
func (s *InMemTranslateStore) translateKey(key string) (_ uint64, err error) {
// Return id if it has been added.
if id, ok := s.lookup[key]; ok {
return id
if id, ok := s.idsByKey[key]; ok {
return id, nil
} else if s.readOnly {
return 0, nil
}
// Generate a new id and update db.
id := uint64(len(s.keys) + 1)
var id uint64
if s.field == "" {
id = GenerateNextPartitionedID(s.index, s.maxID, s.partitionID, s.partitionN)
} else {
id = s.maxID + 1
}
s.set(id, key)
return id
return id, nil
}
// TranslateID converts an integer ID to a string key.
@ -301,10 +344,7 @@ func (s *InMemTranslateStore) TranslateIDs(ids []uint64) ([]string, error) {
}
func (s *InMemTranslateStore) translateID(id uint64) string {
if id == 0 || id > uint64(len(s.keys)) {
return ""
}
return s.keys[id-1]
return s.keysByID[id]
}
// ForceSet writes the id/key pair to the db. Used by replication.
@ -317,8 +357,11 @@ func (s *InMemTranslateStore) ForceSet(id uint64, key string) error {
// set assigns the id/key pair to the store.
func (s *InMemTranslateStore) set(id uint64, key string) {
s.keys = append(s.keys, key)
s.lookup[key] = id
s.keysByID[id] = key
s.idsByKey[key] = id
if id > s.maxID {
s.maxID = id
}
s.notifyWrite()
}
@ -347,7 +390,7 @@ func (s *InMemTranslateStore) EntryReader(ctx context.Context, offset uint64) (T
func (s *InMemTranslateStore) MaxID() (uint64, error) {
s.mu.RLock()
defer s.mu.RUnlock()
return uint64(len(s.keys)), nil
return s.maxID, nil
}
// inMemEntryReader represents a stream of translation entries for an inmem translation store.

View file

@ -26,7 +26,7 @@ import (
)
func TestInMemTranslateStore_TranslateKey(t *testing.T) {
s := pilosa.NewInMemTranslateStore("IDX", "FLD")
s := pilosa.NewInMemTranslateStore("IDX", "FLD", 0, pilosa.DefaultPartitionN)
// Ensure initial key translates to ID 1.
if id, err := s.TranslateKey("foo"); err != nil {
@ -51,7 +51,7 @@ func TestInMemTranslateStore_TranslateKey(t *testing.T) {
}
func TestInMemTranslateStore_TranslateID(t *testing.T) {
s := pilosa.NewInMemTranslateStore("IDX", "FLD")
s := pilosa.NewInMemTranslateStore("IDX", "FLD", 0, pilosa.DefaultPartitionN)
// Setup initial keys.
if _, err := s.TranslateKey("foo"); err != nil {

View file

@ -218,7 +218,7 @@ func (t *ClusterCluster) addCluster(i int, saveTopology bool) (*cluster, error)
}
// holder
h := NewHolder()
h := NewHolder(DefaultPartitionN)
h.Path = path
// cluster