fix id generation

This commit is contained in:
Ben Johnson 2019-12-21 12:03:32 -07:00
parent 82910911dd
commit e3606d6615
8 changed files with 62 additions and 114 deletions

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

@ -169,9 +169,8 @@ func (s *TranslateStore) TranslateKey(key string) (id uint64, _ error) {
return nil
}
if id, err := pilosa.GenerateNextPartitionedID(maxID(tx), s.partitionID, s.partitionN); 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
@ -233,9 +232,8 @@ func (s *TranslateStore) TranslateKeys(keys []string) (ids []uint64, _ error) {
continue
}
if ids[i], err = pilosa.GenerateNextPartitionedID(maxID(tx), s.partitionID, s.partitionN); 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

View file

@ -855,6 +855,10 @@ func (c *cluster) fragSources(to *cluster, idx *Index) (map[string][]*ResizeSour
// 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)
@ -862,21 +866,25 @@ func (c *cluster) shardPartition(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 shard 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(c.partitionN))
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 c.shardPartition(index, (id/ShardWidth)%uint64(c.partitionN))
return shardPartition(index, id/ShardWidth, c.partitionN)
}
// ShardNodes returns a list of nodes that own a fragment. Safe for concurrent use.
@ -2115,7 +2123,6 @@ func (c *cluster) translateIndexKeySet(ctx context.Context, indexName string, ke
g.Go(func() (err error) {
var ids []uint64
println("dbg/translateIndexKeySet", partitionID, len(keys), c.ownsPartition(c.Node.ID, partitionID))
if c.ownsPartition(c.Node.ID, partitionID) {
if ids, err = idx.TranslateStore(partitionID).TranslateKeys(keys); err != nil {
return err
@ -2183,7 +2190,6 @@ func (c *cluster) translateIndexIDSet(ctx context.Context, indexName string, idS
g.Go(func() (err error) {
var keys []string
println("dbg/translateIndexIDSet", partitionID, len(ids), c.ownsPartition(c.Node.ID, partitionID))
if c.ownsPartition(c.Node.ID, partitionID) {
if keys, err = index.TranslateStore(partitionID).TranslateIDs(ids); err != nil {
return err
@ -2198,7 +2204,6 @@ func (c *cluster) translateIndexIDSet(ctx context.Context, indexName string, idS
mu.Lock()
defer mu.Unlock()
for i := range ids {
println("dbg/>>", ids[i], keys[i])
idMap[ids[i]] = keys[i]
}
return nil

View file

@ -3767,7 +3767,6 @@ func (e *executor) translateGroupByCall(index string, idx *Index, c *pql.Call, k
}
func (e *executor) translateResults(ctx context.Context, index string, idx *Index, calls []*pql.Call, results []interface{}) (err error) {
println("dbg/")
span, _ := tracing.StartSpanFromContext(ctx, "Executor.translateResults")
defer span.Finish()
@ -3786,7 +3785,6 @@ func (e *executor) translateResults(ctx context.Context, index string, idx *Inde
}
for i := range results {
println("dbg/results.a", i)
results[i], err = e.translateResult(index, idx, calls[i], results[i], idMap)
if err != nil {
return err

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

View file

@ -725,6 +725,11 @@ func (f *Field) Close() error {
_ = f.rowAttrStore.Close()
}
// Close field translation store.
if err := f.translateStore.Close(); err != nil {
return err
}
// Close all views.
for _, view := range f.viewMap {
if err := view.close(); err != nil {

View file

@ -291,6 +291,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 {

View file

@ -16,7 +16,6 @@ package pilosa
import (
"context"
"fmt"
"io"
"sync"
@ -69,23 +68,15 @@ type TranslateStore interface {
// OpenTranslateStoreFunc represents a function for instantiating and opening a TranslateStore.
type OpenTranslateStoreFunc func(path, index, field string, partitionID, partitionN int) (TranslateStore, error)
// GenerateNextPartitionedID returns the next ID within the same partition.
func GenerateNextPartitionedID(prev uint64, partitionID, partitionN int) (uint64, error) {
// Generate the first available ID for partition if none previously existed.
var id uint64
if prev == 0 {
return (uint64(partitionID) * ShardWidth), nil
} else if (prev/ShardWidth)%uint64(partitionN) != uint64(partitionID) {
return 0, fmt.Errorf("partition id mismatch: id=%d partition=%d", prev, partitionID)
}
// 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.
if id++; (id/ShardWidth)%uint64(partitionN) == uint64(partitionID) {
return id, nil
// 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
}
}
// Otherwise jump to the next possible shard that is in the same partition.
return (id - ShardWidth) + (ShardWidth * uint64(partitionN)), nil
}
// TranslateEntryReader represents a stream of translation entries.
@ -330,9 +321,7 @@ func (s *InMemTranslateStore) translateKey(key string) (_ uint64, err error) {
// Generate a new id and update db.
var id uint64
if s.field == "" {
if id, err = GenerateNextPartitionedID(s.maxID, s.partitionID, s.partitionN); err != nil {
return 0, err
}
id = GenerateNextPartitionedID(s.index, s.maxID, s.partitionID, s.partitionN)
} else {
id = s.maxID + 1
}