diff --git a/api_test.go b/api_test.go index b6cfa8d07..aabbaa700 100644 --- a/api_test.go +++ b/api_test.go @@ -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{ diff --git a/boltdb/translate.go b/boltdb/translate.go index 4168b1e14..dbdf5a867 100644 --- a/boltdb/translate.go +++ b/boltdb/translate.go @@ -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 diff --git a/cluster.go b/cluster.go index 05ba2a27c..fdbf771ea 100644 --- a/cluster.go +++ b/cluster.go @@ -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 diff --git a/executor.go b/executor.go index a5fbb14d0..3f9a48d67 100644 --- a/executor.go +++ b/executor.go @@ -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 diff --git a/executor_test.go b/executor_test.go index ec8865c8a..c2e679cbc 100644 --- a/executor_test.go +++ b/executor_test.go @@ -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) } } } diff --git a/field.go b/field.go index efd84933e..3fea5440d 100644 --- a/field.go +++ b/field.go @@ -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 { diff --git a/index.go b/index.go index bcf63f1f1..a8adfc674 100644 --- a/index.go +++ b/index.go @@ -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 { diff --git a/translate.go b/translate.go index 1e3e2094c..158677051 100644 --- a/translate.go +++ b/translate.go @@ -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 }