Merge branch 'master' into grpc-crd-sql

This commit is contained in:
Cody Soyland 2020-08-12 14:21:53 -05:00 • committed by GitHub
commit 21456827ae
No known key found for this signature in database
GPG key ID: 4AEE18F83AFDEB23
20 changed files with 4060 additions and 362 deletions

27
api.go
View file

@ -1066,10 +1066,6 @@ func (api *API) Import(ctx context.Context, req *ImportRequest, opts ...ImportOp
return errors.Wrap(err, "getting index and field")
}
// Obtain transaction.
tx := index.Txf.NewTx(Txo{Write: true, Index: index})
defer tx.Rollback()
if err := req.ValidateWithTimestamp(index.CreatedAt(), field.CreatedAt()); err != nil {
return errors.Wrap(err, "validating import value request")
}
@ -1162,12 +1158,23 @@ func (api *API) Import(ctx context.Context, req *ImportRequest, opts ...ImportOp
// Import columnIDs into existence field.
if !options.Clear {
if err := importExistenceColumns(tx, index, req.ColumnIDs); err != nil {
api.server.logger.Printf("import existence error: index=%s, field=%s, shard=%d, columns=%d, err=%s", req.Index, req.Field, req.Shard, len(req.ColumnIDs), err)
if err := func() error {
tx := index.Txf.NewTx(Txo{Write: true, Index: index})
defer tx.Rollback()
if err := importExistenceColumns(tx, index, req.ColumnIDs); err != nil {
api.server.logger.Printf("import existence error: index=%s, field=%s, shard=%d, columns=%d, err=%s", req.Index, req.Field, req.Shard, len(req.ColumnIDs), err)
return err
}
return tx.Commit()
}(); err != nil {
return errors.Wrap(err, "importing existence columns")
}
}
tx := index.Txf.NewTx(Txo{Write: true, Index: index})
defer tx.Rollback()
// Import into fragment.
err = field.Import(tx, req.RowIDs, req.ColumnIDs, timestamps, opts...)
if err != nil {
@ -1198,10 +1205,6 @@ func (api *API) ImportValue(ctx context.Context, req *ImportValueRequest, opts .
return errors.Wrap(err, "validating import value request")
}
// Obtain transaction.
tx := index.Txf.NewTx(Txo{Write: true, Index: index})
defer tx.Rollback()
// Set up import options.
options, err := setUpImportOptions(opts...)
if err != nil {
@ -1257,6 +1260,10 @@ func (api *API) ImportValue(ctx context.Context, req *ImportValueRequest, opts .
// if we're importing into a specific shard
if req.Shard != math.MaxUint64 {
// Obtain transaction.
tx := index.Txf.NewTx(Txo{Write: true, Index: index})
defer tx.Rollback()
// Check that column IDs match the stated shard.
if s1, s2 := req.ColumnIDs[0]/ShardWidth, req.ColumnIDs[len(req.ColumnIDs)-1]/ShardWidth; s1 != s2 && s2 != req.Shard {
return errors.Errorf("shard %d specified, but import spans shards %d to %d", req.Shard, s1, s2)

View file

@ -165,8 +165,6 @@ func TestAPI_ImportColumnAttrs(t *testing.T) {
}
func TestAPI_Import(t *testing.T) {
skipForRBF(t)
c := test.MustRunCluster(t, 2,
[]server.CommandOption{
server.OptCommandServerOptions(
@ -278,8 +276,6 @@ func TestAPI_Import(t *testing.T) {
}
func TestAPI_ImportValue(t *testing.T) {
skipForRBF(t)
c := test.MustRunCluster(t, 2,
[]server.CommandOption{
server.OptCommandServerOptions(

View file

@ -1590,8 +1590,11 @@ func (c *cluster) followResizeInstruction(instr *ResizeInstruction) error {
}
// Create view.
v, err := f.createViewIfNotExists(src.View)
if err != nil {
var v *view
if err := func() (err error) {
v, err = f.createViewIfNotExists(src.View)
return err
}(); err != nil {
return errors.Wrap(err, "creating view")
}

View file

@ -15,7 +15,7 @@ Is equivalent to `GET /schema` and returns the same response.
### List index schema
`GET /index/<index-name>`
`GET /index/{index-name}`
Returns the schema of the specified index in JSON.
@ -48,7 +48,7 @@ curl -XGET localhost:10101/index/user
### Create index
`POST /index/<index-name>`
`POST /index/{index-name}`
Creates an index with the given name.
@ -79,7 +79,7 @@ curl -XDELETE localhost:10101/index/user
### Query index
`POST /index/<index-name>/query`
`POST /index/{index-name}/query`
Sends a [query](../query-language/) to the Pilosa server with the given index. The request body is UTF-8 encoded text and response body is in JSON by default.
@ -137,7 +137,7 @@ By default, all bits and attributes (*for `Row` queries only*) are returned. In
### Import Data
`POST /index/<index-name>/field/<field-name>/import`
`POST /index/{index-name}/field/{field-name}/import`
Supports high-rate data ingest to a particular shard of a particular field. The
official client libraries use this endpoint for their import functionality - it
@ -185,7 +185,7 @@ message ImportRequest {
### Create field
`POST /index/<index-name>/field/<field-name>`
`POST /index/{index-name}/field/{field-name}`
Creates a field in the given index with the given name.
@ -241,7 +241,7 @@ curl localhost:10101/index/repository/field/stats \
### Remove field
`DELETE /index/<index-name>/field/<field-name>`
`DELETE /index/{index-name}/field/{field-name}`
Removes the given field.

View file

@ -58,7 +58,7 @@ nav = []
<strong id="range-bsi">[Row (BSI)](../query-language/#row-bsi):</strong> A [PQL](#pql) query that returns bits based on comparison to integers stored in [BSI](#bsi) [fields](#field).
<strong id="rows">[Rows](../query-language/#rows):<strong> A [PQL](#pql) query that returns a list of row IDs in the given field which have at least one bit set. The field argument is mandatory, the others are optional. `Rows` is the primary argument used with the [GroupBy](#groupby) query.
<strong id="rows">[Rows](../query-language/#rows):</strong> A [PQL](#pql) query that returns a list of row IDs in the given field which have at least one bit set. The field argument is mandatory, the others are optional. `Rows` is the primary argument used with the [GroupBy](#groupby) query.
<strong id="slice">[Slice](../data-model/#shard):</strong> Prior to Pilosa 1.0, shards were known as slices.

View file

@ -429,9 +429,9 @@ Union([ROW_CALL ...])
**Description:**
Union performs a logical OR on the results of all `ROW_CALL` queries passed to it.
Union performs a set union on the column indexes in the results of all `ROW_CALL` queries passed to it. In comparison to a relational query, this is similar to combining clauses in the "OR" sense.
**Result Type:** object with attrs and bits
**Result Type:** object with attrs and columns
attrs will always be empty
@ -457,7 +457,7 @@ Intersect(<ROW_CALL>, [ROW_CALL ...])
**Description:**
Intersect performs a logical AND on the results of all `ROW_CALL` queries passed to it.
Intersect performs a set intersection on the column indexes in the results of all `ROW_CALL` queries passed to it. In comparison to a relational query, this is similar to combining clauses in the "AND" sense.
**Result Type:** object with attrs and columns
@ -571,6 +571,38 @@ Not(Row(stargazer=1))
* columns are repositories that were not starred by user 1
#### Limit
**Spec:**
```
Limit(<ROW_CALL>, [limit=<UINT>], [offset=<UINT>])
```
**Description:**
Limit executes a `ROW_CALL` and returns a subset of the results.
If a limit of `n` is specified, then this query will return the first `n` results of the row call.
If an offset of `m` is specified, then this query will skip the first `m` results of the row call.
If both a limit and offset are specified, the offset is applied before the limit.
This can be used to implement pagination.
**Result Type:** object with attrs and columns
attrs will always be empty
**Examples:**
Find the second column that has a bit set in the given row.
```request
Limit(Row(stargazer=1), limit=1, offset=1)
```
```response
{"results":[{"attrs":{},"columns":[30]}]}
```
* columns are repositories that were starred by user 1
#### Count
**Spec:**
@ -792,6 +824,33 @@ Options(Row(f1=10), shards=[0, 2])
{"attrs":{},"columns":[100, 2097152]}
```
#### Row Constant
**Spec:**
```
ConstRow(columns=<[]COLUMN>)
```
**Description:**
`ConstRow` provides a constant bitmap value that can be used in place of a `Row` call.
The columns can be specified as integer IDs or strings.
**Result Type:** row value columns.
e.g. `{"attrs":{},"columns":[10, 20]}`
**Examples:**
Filter specified columns to only those with a bit set in row 1 of the field `stargazer` (repositories that are starred by user 1):
```request
Intersect(ConstRow(columns=[10, 20, 30]), Row(stargazer=1))
```
```response
{"attrs":{},"columns":[10, 20]}
```
#### Rows
**Spec:**
@ -850,6 +909,41 @@ Rows(job, like="%t")
{"rows":null,"keys":["management","student"]}
```
#### Extract
**Spec:**
```
Extract(<ROW_CALL>, [<ROWS_CALL>...])
```
**Description:**
Extract intersects a set of columns with a set of rows in order to extract a subset of the index.
The result is a table consisting of the matched columns and the rows which they intersect.
This is similar to a select query in a SQL database.
**Result Type:** Object with an array of the selected fields and an array of the selected columns.
The column array contains objects containing a column identifier and an array of field values.
Field values are typed as such:
- Bool Field - boolean or null
- Mutex Field (unkeyed) - 64-bit unsigned integer or null
- Mutex Field (keyed) - string or null
- Integer Field - 64-bit signed integer or null
- Decimal Field - Pilosa decimal value or null
- Set Field (unkeyed) - array of 64-bit unsigned integers
- Set Field (keyed) - array of strings
- Time Field - same as the equivalent Set
**Examples:**
List all stargazers who have starred repository 1, and the full set of repositories they have starred:
```request
Extract(Row(stargazer=1), Rows(stargazer))
```
```response
{"fields":[{"name":"stargazer","type":"set"}],"columns":[{"column":3,"rows":[[1, 2, 3]]}]}
```
#### Group By
**Spec:**
@ -949,7 +1043,7 @@ UnionRows([ROWSET_CALL ...])
UnionRows performs a logical OR on the rows matched by the results of all `ROWSET_CALL` queries passed to it.
**Result Type:** object with attrs and bits
**Result Type:** object with attrs and columns
attrs will always be empty
@ -963,4 +1057,4 @@ UnionRows(Rows(stargazer))
{"attrs":{},"columns":[10, 20, 30]}
```
* columns are repositories that were starred by any user
* columns are repositories that were starred by any user

View file

@ -512,12 +512,18 @@ func (s Serializer) encodeQueryResponse(m *pilosa.QueryResponse) *internal.Query
case pilosa.RowIDs:
pb.Results[i].Type = queryResultTypeRowIDs
pb.Results[i].RowIDs = result
case pilosa.ExtractedIDMatrix:
pb.Results[i].Type = queryResultTypeExtractedIDMatrix
pb.Results[i].ExtractedIDMatrix = s.endcodeExtractedIDMatrix(result)
case []pilosa.GroupCount:
pb.Results[i].Type = queryResultTypeGroupCounts
pb.Results[i].GroupCounts = s.encodeGroupCounts(result)
case pilosa.RowIdentifiers:
pb.Results[i].Type = queryResultTypeRowIdentifiers
pb.Results[i].RowIdentifiers = s.encodeRowIdentifiers(result)
case pilosa.ExtractedTable:
pb.Results[i].Type = queryResultTypeExtractedTable
pb.Results[i].ExtractedTable = s.encodeExtractedTable(result)
case pilosa.Pair:
pb.Results[i].Type = queryResultTypePair
pb.Results[i].Pairs = []*internal.Pair{s.encodePair(result)}
@ -1313,6 +1319,8 @@ const (
queryResultTypePair
queryResultTypePairField
queryResultTypeSignedRow
queryResultTypeExtractedIDMatrix
queryResultTypeExtractedTable
)
func (s Serializer) decodeQueryResult(pb *internal.QueryResult) interface{} {
@ -1343,6 +1351,10 @@ func (s Serializer) decodeQueryResult(pb *internal.QueryResult) interface{} {
return s.decodePair(pb.Pairs[0])
case queryResultTypePairField:
return s.decodePairField(pb.PairField)
case queryResultTypeExtractedIDMatrix:
return s.decodeExtractedIDMatrix(pb.ExtractedIDMatrix)
case queryResultTypeExtractedTable:
return s.decodeExtractedTable(pb.ExtractedTable)
}
panic(fmt.Sprintf("unknown type: %d", pb.Type))
}
@ -1410,6 +1422,83 @@ func (s Serializer) decodeAttr(attr *internal.Attr) (key string, value interface
}
}
func (s Serializer) decodeExtractedIDMatrix(m *internal.ExtractedIDMatrix) pilosa.ExtractedIDMatrix {
cols := make([]pilosa.ExtractedIDColumn, len(m.Columns))
for i, c := range m.Columns {
rows := make([][]uint64, len(c.Vals))
for j, r := range c.Vals {
rows[j] = r.IDs
}
cols[i] = pilosa.ExtractedIDColumn{
ColumnID: c.ID,
Rows: rows,
}
}
return pilosa.ExtractedIDMatrix{
Fields: m.Fields,
Columns: cols,
}
}
func (s Serializer) decodeExtractedTable(t *internal.ExtractedTable) pilosa.ExtractedTable {
fields := make([]pilosa.ExtractedTableField, len(t.Fields))
for i, f := range t.Fields {
fields[i] = pilosa.ExtractedTableField{
Name: f.Name,
Type: f.Type,
}
}
columns := make([]pilosa.ExtractedTableColumn, len(t.Columns))
for i, c := range t.Columns {
var col pilosa.KeyOrID
switch kid := c.KeyOrID.(type) {
case *internal.ExtractedTableColumn_ID:
col = pilosa.KeyOrID{
ID: kid.ID,
}
case *internal.ExtractedTableColumn_Key:
col = pilosa.KeyOrID{
Keyed: true,
Key: kid.Key,
}
}
rows := make([]interface{}, len(c.Values))
for j, v := range rows {
var val interface{}
switch v := v.(type) {
case *internal.ExtractedTableValue_IDs:
val = v.IDs.IDs
case *internal.ExtractedTableValue_Keys:
val = v.Keys.Keys
case *internal.ExtractedTableValue_BSIValue:
val = v.BSIValue
case *internal.ExtractedTableValue_MutexID:
val = v.MutexID
case *internal.ExtractedTableValue_MutexKey:
val = v.MutexKey
case *internal.ExtractedTableValue_Bool:
val = v.Bool
}
rows[j] = val
}
columns[i] = pilosa.ExtractedTableColumn{
Column: col,
Rows: rows,
}
}
return pilosa.ExtractedTable{
Fields: fields,
Columns: columns,
}
}
func (s Serializer) decodeRowIdentifiers(a *internal.RowIdentifiers) *pilosa.RowIdentifiers {
return &pilosa.RowIdentifiers{
Rows: a.Rows,
@ -1580,6 +1669,98 @@ func (s Serializer) encodeFieldRows(a []pilosa.FieldRow) []*internal.FieldRow {
return other
}
func (s Serializer) endcodeExtractedIDMatrix(m pilosa.ExtractedIDMatrix) *internal.ExtractedIDMatrix {
cols := make([]*internal.ExtractedIDColumn, len(m.Columns))
for i, v := range m.Columns {
vals := make([]*internal.IDList, len(v.Rows))
for j, f := range v.Rows {
vals[j] = &internal.IDList{IDs: f}
}
cols[i] = &internal.ExtractedIDColumn{
ID: v.ColumnID,
Vals: vals,
}
}
return &internal.ExtractedIDMatrix{
Fields: m.Fields,
Columns: cols,
}
}
func (s Serializer) encodeExtractedTable(t pilosa.ExtractedTable) *internal.ExtractedTable {
fields := make([]*internal.ExtractedTableField, len(t.Fields))
for i, f := range t.Fields {
fields[i] = &internal.ExtractedTableField{
Name: f.Name,
Type: f.Type,
}
}
cols := make([]*internal.ExtractedTableColumn, len(t.Columns))
for i, c := range t.Columns {
var col internal.ExtractedTableColumn
if c.Column.Keyed {
col.KeyOrID = &internal.ExtractedTableColumn_Key{Key: c.Column.Key}
} else {
col.KeyOrID = &internal.ExtractedTableColumn_ID{ID: c.Column.ID}
}
rows := make([]*internal.ExtractedTableValue, len(c.Rows))
for j, v := range c.Rows {
switch v := v.(type) {
case []uint64:
rows[j] = &internal.ExtractedTableValue{
Value: &internal.ExtractedTableValue_IDs{
IDs: &internal.IDList{
IDs: v,
},
},
}
case []string:
rows[j] = &internal.ExtractedTableValue{
Value: &internal.ExtractedTableValue_Keys{
Keys: &internal.KeyList{
Keys: v,
},
},
}
case int64:
rows[j] = &internal.ExtractedTableValue{
Value: &internal.ExtractedTableValue_BSIValue{
BSIValue: v,
},
}
case uint64:
rows[j] = &internal.ExtractedTableValue{
Value: &internal.ExtractedTableValue_MutexID{
MutexID: v,
},
}
case string:
rows[j] = &internal.ExtractedTableValue{
Value: &internal.ExtractedTableValue_MutexKey{
MutexKey: v,
},
}
case bool:
rows[j] = &internal.ExtractedTableValue{
Value: &internal.ExtractedTableValue_Bool{
Bool: v,
},
}
}
}
col.Values = rows
cols[i] = &col
}
return &internal.ExtractedTable{
Fields: fields,
Columns: cols,
}
}
func (s Serializer) encodePairs(a pilosa.Pairs) []*internal.Pair {
other := make([]*internal.Pair, len(a))
for i := range a {

View file

@ -289,6 +289,10 @@ func (e *executor) safeCopy(resp QueryResponse) (out QueryResponse) {
out.Results = append(out.Results, x)
case []GroupCount:
out.Results = append(out.Results, x)
case ExtractedTable:
out.Results = append(out.Results, x)
case ExtractedIDMatrix:
out.Results = append(out.Results, x)
case RowIdentifiers:
// no bitmap material, so should be ok to skip Clone()
out.Results = append(out.Results, x)
@ -519,7 +523,6 @@ func (e *executor) execute(ctx context.Context, tx Tx, index string, q *pql.Quer
}
// preprocessQuery expands any calls that need preprocessing.
// So far, this only needs to process UnionRows.
func (e *executor) preprocessQuery(ctx context.Context, tx Tx, index string, c *pql.Call, shards []uint64, opt *execOptions) (*pql.Call, error) {
switch c.Name {
case "UnionRows":
@ -601,6 +604,88 @@ func (e *executor) preprocessQuery(ctx context.Context, tx Tx, index string, c *
Children: rows,
}, nil
case "ConstRow":
// Fetch user-provided columns list.
cols, _ := c.Args["columns"].([]interface{})
var ids []uint64
var keys []string
for _, c := range cols {
switch c := c.(type) {
case uint64:
ids = append(ids, c)
case int64:
ids = append(ids, uint64(c))
case string:
keys = append(keys, c)
default:
return nil, errors.Errorf("invalid column identifier %v of type %T", c, c)
}
}
// Translate keys to IDs.
if len(keys) > 0 {
keyIDs, err := e.Cluster.translateIndexKeys(ctx, index, keys)
if err != nil {
return nil, errors.Wrap(err, "translating column IDs in ConstRow")
}
ids = append(ids, keyIDs...)
}
// Split IDs by shard.
shardSet := make(map[uint64][]uint64)
for _, id := range ids {
shardSet[id/ShardWidth] = append(shardSet[id/ShardWidth], id)
}
// Convert ID sets to per-shard Row objects.
precomputed := make(map[uint64]interface{})
for _, s := range shards {
precomputed[s] = NewRow(shardSet[s]...)
}
// Generate a precomputed call with the data.
return &pql.Call{
Name: "Precomputed",
Precomputed: precomputed,
}, nil
case "All":
_, hasLimit, err := c.UintArg("limit")
if err != nil {
return nil, err
}
_, hasOffset, err := c.UintArg("offset")
if err != nil {
return nil, err
}
if !hasLimit && !hasOffset {
return c, nil
}
// Rewrite the All() w/ limit to Limit(All()).
c.Children = []*pql.Call{
{
Name: "All",
},
}
c.Name = "Limit"
fallthrough
case "Limit":
if len(c.Children) != 1 {
return nil, errors.Errorf("expected 1 child of limit call but got %d", len(c.Children))
}
res, err := e.preprocessQuery(ctx, tx, index, c.Children[0], shards, opt)
if err != nil {
return nil, err
}
c.Children[0] = res
err = e.executeLimitCall(ctx, tx, index, c, shards, opt)
if err != nil {
return nil, err
}
return c, nil
default:
// Recurse through child calls.
out := make([]*pql.Call, len(c.Children))
@ -714,6 +799,9 @@ func (e *executor) executeCall(ctx context.Context, tx Tx, index string, c *pql.
case "Rows":
statFn()
return e.executeRows(ctx, tx, index, c, shards, opt)
case "Extract":
statFn()
return e.executeExtract(ctx, tx, index, c, shards, opt)
case "GroupBy":
statFn()
return e.executeGroupBy(ctx, tx, index, c, shards, opt)
@ -725,9 +813,6 @@ func (e *executor) executeCall(ctx context.Context, tx Tx, index string, c *pql.
case "FieldValue":
statFn()
return e.executeFieldValueCall(ctx, tx, index, c, shards, opt)
case "All":
statFn()
return e.executeAllCall(ctx, tx, index, c, shards, opt)
case "Precomputed":
return e.executePrecomputedCall(ctx, tx, index, c, shards, opt)
default:
@ -925,25 +1010,20 @@ func (e *executor) executeFieldValueCallShard(ctx context.Context, tx Tx, field
return other, nil
}
// executeAllCall executes an All() call.
func (e *executor) executeAllCall(ctx context.Context, tx Tx, index string, c *pql.Call, shards []uint64, opt *execOptions) (*Row, error) {
rslt := NewRow()
// executeLimitCall executes a Limit() call, **rewriting it to a precomputed call**.
func (e *executor) executeLimitCall(ctx context.Context, tx Tx, index string, c *pql.Call, shards []uint64, opt *execOptions) error {
bitmapCall := c.Children[0]
var limit uint64
var offset uint64
if lim, hasLimit, err := c.UintArg("limit"); err != nil {
return nil, errors.Wrap(err, "getting limit")
} else if hasLimit && lim > 0 {
limit = uint64(lim)
limit, hasLimit, err := c.UintArg("limit")
if err != nil {
return errors.Wrap(err, "getting limit")
}
if off, hasOffset, err := c.UintArg("offset"); err != nil {
return nil, errors.Wrap(err, "getting offset")
} else if hasOffset && off > 0 {
offset = uint64(off)
offset, _, err := c.UintArg("offset")
if err != nil {
return errors.Wrap(err, "getting offset")
}
if limit == 0 {
if !hasLimit {
limit = math.MaxUint64
}
@ -954,12 +1034,34 @@ func (e *executor) executeAllCall(ctx context.Context, tx Tx, index string, c *p
// got tracks the number of records gotten to that point.
var got uint64
c.Precomputed = make(map[uint64]interface{})
for _, shard := range shards {
row, err := e.executeAllCallMapReduce(ctx, tx, index, c, shard, opt)
if err != nil {
return nil, errors.Wrap(err, "executing map reduce on shard")
// Execute calls in bulk on each remote node and merge.
mapFn := func(ctx context.Context, shard uint64) (interface{}, error) {
return e.executeBitmapCallShard(ctx, tx, index, bitmapCall, shard)
}
// Merge returned results at coordinating node.
reduceFn := func(ctx context.Context, prev, v interface{}) interface{} {
if err := ctx.Err(); err != nil {
return err
}
other, _ := prev.(*Row)
if other == nil {
other = NewRow()
}
other.Merge(v.(*Row))
return other
}
result, err := e.mapReduce(ctx, index, []uint64{shard}, c, opt, mapFn, reduceFn)
if err != nil {
return errors.Wrap(err, "limit map reduce")
}
row, _ := result.(*Row)
segCnt := row.Count()
// If this segment doesn't reach the offset, skip it.
@ -972,14 +1074,14 @@ func (e *executor) executeAllCall(ctx context.Context, tx Tx, index string, c *p
// (or it has exactly enough).
if segCnt-skip <= limit-got {
if skip == 0 {
rslt.Merge(row)
c.Precomputed[shard] = row
} else {
cols := row.Columns()
partialRow := NewRow()
for _, bit := range cols[skip:] {
partialRow.SetBit(bit)
}
rslt.Merge(partialRow)
c.Precomputed[shard] = partialRow
}
got += segCnt - skip
// In the case where this segment exactly fulfills the limit, break.
@ -996,42 +1098,12 @@ func (e *executor) executeAllCall(ctx context.Context, tx Tx, index string, c *p
for _, bit := range cols[skip : skip+limit-got] {
partialRow.SetBit(bit)
}
rslt.Merge(partialRow)
c.Precomputed[shard] = partialRow
break
}
return rslt, nil
}
// executeAllCallMapReduce executes a single shard of the All() call
// using the executor.mapReduce() method.
func (e *executor) executeAllCallMapReduce(ctx context.Context, tx Tx, index string, c *pql.Call, shard uint64, opt *execOptions) (*Row, error) {
// Execute calls in bulk on each remote node and merge.
mapFn := func(ctx context.Context, shard uint64) (interface{}, error) {
return e.executeAllCallShard(ctx, tx, index, c, shard)
}
// Merge returned results at coordinating node.
reduceFn := func(ctx context.Context, prev, v interface{}) interface{} {
if err := ctx.Err(); err != nil {
return err
}
other, _ := prev.(*Row)
if other == nil {
other = NewRow()
}
other.Merge(v.(*Row))
return other
}
result, err := e.mapReduce(ctx, index, []uint64{shard}, c, opt, mapFn, reduceFn)
if err != nil {
return nil, errors.Wrap(err, "map reduce")
}
row, _ := result.(*Row)
return row, nil
c.Name = "Precomputed"
return nil
}
// executeIncludesColumnCallShard
@ -1407,7 +1479,7 @@ func (e *executor) executeBitmapCallShard(ctx context.Context, tx Tx, index stri
return e.executeNotShard(ctx, tx, index, c, shard)
case "Shift":
return e.executeShiftShard(ctx, tx, index, c, shard)
case "All": // Allow a shard computation to use All() (note, limit/offset not applied)
case "All": // Allow a shard computation to use All()
return e.executeAllCallShard(ctx, tx, index, c, shard)
case "Distinct":
return nil, errors.New("Distinct shouldn't be hit as a bitmap call")
@ -2590,8 +2662,11 @@ func (e *executor) executeRowsShard(ctx context.Context, tx Tx, index string, fi
// in order to represent `Rows` for the field.
var views = []string{viewStandard}
// Handle `time` fields.
if f.Type() == FieldTypeTime {
// Handle `int` and `time` fields.
switch f.Type() {
case FieldTypeInt:
return nil, errors.New("int fields not supported by Rows() query")
case FieldTypeTime:
var err error
// Parse "from" time, if set.
@ -2712,6 +2787,298 @@ func (e *executor) executeRowsShard(ctx context.Context, tx Tx, index string, fi
return rowIDs, nil
}
type ExtractedTableField struct {
Name string `json:"name"`
Type string `json:"type"`
}
type KeyOrID struct {
ID uint64
Key string
Keyed bool
}
func (kid KeyOrID) MarshalJSON() ([]byte, error) {
if kid.Keyed {
return json.Marshal(kid.Key)
}
return json.Marshal(kid.ID)
}
type ExtractedTableColumn struct {
Column KeyOrID `json:"column"`
Rows []interface{} `json:"rows"`
}
type ExtractedTable struct {
Fields []ExtractedTableField `json:"fields"`
Columns []ExtractedTableColumn `json:"columns"`
}
type ExtractedIDColumn struct {
ColumnID uint64
Rows [][]uint64
}
type ExtractedIDMatrix struct {
Fields []string
Columns []ExtractedIDColumn
}
func (e *ExtractedIDMatrix) Append(m ExtractedIDMatrix) {
e.Columns = append(e.Columns, m.Columns...)
if e.Fields == nil {
e.Fields = m.Fields
}
}
func (e *executor) executeExtract(ctx context.Context, tx Tx, index string, c *pql.Call, shards []uint64, opt *execOptions) (ExtractedIDMatrix, error) {
// Extract the column filter call.
if len(c.Children) < 1 {
return ExtractedIDMatrix{}, errors.New("missing column filter in Extract")
}
filter := c.Children[0]
// Extract fields from rows calls.
fields := make([]string, len(c.Children)-1)
for i, rows := range c.Children[1:] {
if rows.Name != "Rows" {
return ExtractedIDMatrix{}, errors.Errorf("child call of Extract is %q but expected Rows", rows.Name)
}
var fieldName string
var ok bool
for k, v := range rows.Args {
switch k {
case "field", "_field":
fieldName = v.(string)
ok = true
default:
return ExtractedIDMatrix{}, errors.Errorf("unsupported Rows argument for Extract: %q", k)
}
}
if !ok {
return ExtractedIDMatrix{}, errors.New("missing field specification in Rows")
}
fields[i] = fieldName
}
// Execute calls in bulk on each remote node and merge.
mapFn := func(ctx context.Context, shard uint64) (interface{}, error) {
return e.executeExtractShard(ctx, tx, index, fields, filter, shard)
}
// Merge returned results at coordinating node.
reduceFn := func(ctx context.Context, prev, v interface{}) interface{} {
other, _ := prev.(ExtractedIDMatrix)
if err := ctx.Err(); err != nil {
return err
}
other.Append(v.(ExtractedIDMatrix))
return other
}
// Get full result set.
other, err := e.mapReduce(ctx, index, shards, c, opt, mapFn, reduceFn)
if err != nil {
return ExtractedIDMatrix{}, err
}
results, _ := other.(ExtractedIDMatrix)
sort.Slice(results.Columns, func(i, j int) bool {
return results.Columns[i].ColumnID < results.Columns[j].ColumnID
})
return results, nil
}
func mergeBits(bits *Row, mask uint64, out map[uint64]uint64) {
for _, v := range bits.Columns() {
out[v] |= mask
}
}
var trueRowFakeID = []uint64{1}
var falseRowFakeID = []uint64{0}
func (e *executor) executeExtractShard(ctx context.Context, tx Tx, index string, fields []string, filter *pql.Call, shard uint64) (ExtractedIDMatrix, error) {
// Execute filter.
colsBitmap, err := e.executeBitmapCallShard(ctx, tx, index, filter, shard)
if err != nil {
return ExtractedIDMatrix{}, errors.Wrap(err, "failed to get extraction column filter")
}
// Fetch index.
idx := e.Holder.Index(index)
if idx == nil {
return ExtractedIDMatrix{}, ErrIndexNotFound
}
// Decompress columns bitmap.
cols := colsBitmap.Columns()
// Generate a matrix to stuff the results into.
m := make([]ExtractedIDColumn, len(cols))
{
rowsBuf := make([][]uint64, len(m)*len(fields))
for i, c := range cols {
m[i] = ExtractedIDColumn{
ColumnID: c,
Rows: rowsBuf[i*len(fields) : (i+1)*len(fields) : (i+1)*len(fields)],
}
}
}
if len(m) == 0 {
return ExtractedIDMatrix{
Fields: fields,
Columns: m,
}, nil
}
mLookup := make(map[uint64]int)
for i, j := range cols {
mLookup[j] = i
}
// Process fields.
for i, name := range fields {
// Look up the field.
field := idx.Field(name)
if field == nil {
return ExtractedIDMatrix{}, ErrFieldNotFound
}
switch field.Type() {
case FieldTypeSet, FieldTypeMutex, FieldTypeTime:
// Handle a set field by listing the rows and then intersecting them with the filter.
// Extract the standard view fragment.
fragment := e.Holder.fragment(index, name, viewStandard, shard)
if fragment == nil {
// There is nothing here.
continue
}
// List all rows in the standard view.
rows, err := fragment.rows(ctx, tx, 0)
if err != nil {
return ExtractedIDMatrix{}, errors.Wrap(err, "listing rows in set field")
}
// Loop over each row and scan the intersection with the filter.
for _, rowID := range rows {
// Load row from fragment.
row, err := fragment.row(tx, rowID)
if err != nil {
return ExtractedIDMatrix{}, errors.Wrap(err, "loading row from fragment")
}
// Apply column filter to row.
row = row.Intersect(colsBitmap)
// Rotate vector into the matrix.
for _, columnID := range row.Columns() {
fieldSlot := &m[mLookup[columnID]].Rows[i]
*fieldSlot = append(*fieldSlot, rowID)
}
}
case FieldTypeBool:
// Handle bool fields by scanning the true and false rows and assigning an integer.
// Extract the standard view fragment.
fragment := e.Holder.fragment(index, name, viewStandard, shard)
if fragment == nil {
// There is nothing here.
continue
}
// Fetch true and false rows.
trueRow, err := fragment.row(tx, trueRowID)
if err != nil {
return ExtractedIDMatrix{}, errors.Wrap(err, "loading true row from fragment")
}
falseRow, err := fragment.row(tx, falseRowID)
if err != nil {
return ExtractedIDMatrix{}, errors.Wrap(err, "loading true row from fragment")
}
// Fetch values by column.
for j := range m {
col := m[j].ColumnID
switch {
case trueRow.Includes(col):
m[j].Rows[i] = trueRowFakeID
case falseRow.Includes(col):
m[j].Rows[i] = falseRowFakeID
}
}
case FieldTypeInt, FieldTypeDecimal:
// Handle an int/decimal field by rotating a BSI matrix.
// Extract the BSI view fragment.
fragment := e.Holder.fragment(index, name, viewBSIGroupPrefix+name, shard)
if fragment == nil {
// There is nothing here.
continue
}
// Load the BSI group.
bsig := field.bsiGroup(name)
if bsig == nil {
return ExtractedIDMatrix{}, ErrBSIGroupNotFound
}
// Load the BSI exists bit.
exists, err := fragment.row(tx, bsiExistsBit)
if err != nil {
return ExtractedIDMatrix{}, errors.Wrap(err, "loading BSI exists bit from fragment")
}
// Filter BSI exists bit by selected columns.
exists = exists.Intersect(colsBitmap)
if !exists.Any() {
// No relevant BSI values are present in this fragment.
continue
}
// Populate a map with the BSI data.
data := make(map[uint64]uint64)
mergeBits(exists, 0, data)
// Copy in the sign bit.
sign, err := fragment.row(tx, bsiSignBit)
if err != nil {
return ExtractedIDMatrix{}, errors.Wrap(err, "loading BSI sign bit from fragment")
}
sign = sign.Intersect(exists)
mergeBits(sign, 1<<63, data)
// Copy in the significand.
for i := uint(0); i < bsig.BitDepth; i++ {
bits, err := fragment.row(tx, bsiOffsetBit+uint64(i))
if err != nil {
return ExtractedIDMatrix{}, errors.Wrap(err, "loading BSI significand bit from fragment")
}
bits = bits.Intersect(exists)
mergeBits(bits, 1<<i, data)
}
// Store the results back into the matrix.
for columnID, val := range data {
// Convert to two's complement.
val = uint64((2*(int64(val)>>63) + 1) * int64(val&^(1<<63)))
m[mLookup[columnID]].Rows[i] = []uint64{val}
}
}
}
// Emit the final matrix.
// Like RowIDs, this is an internal type and will need to be converted.
return ExtractedIDMatrix{
Fields: fields,
Columns: m,
}, nil
}
func (e *executor) executeRowShard(ctx context.Context, tx Tx, index string, c *pql.Call, shard uint64) (*Row, error) {
span, _ := tracing.StartSpanFromContext(ctx, "Executor.executeRowShard")
@ -4452,18 +4819,24 @@ func (e *executor) translateResults(ctx context.Context, index string, idx *Inde
}
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
}
switch result := result.(type) {
case *Row:
if !idx.Keys() {
return nil
}
for _, segment := range row.Segments() {
for _, col := range segment.Columns() {
idSet[col] = struct{}{}
for _, segment := range result.Segments() {
for _, col := range segment.Columns() {
idSet[col] = struct{}{}
}
}
case ExtractedIDMatrix:
for _, col := range result.Columns {
idSet[col.ColumnID] = struct{}{}
}
}
return nil
}
@ -4630,6 +5003,151 @@ func (e *executor) translateResult(ctx context.Context, index string, idx *Index
}
return other, nil
case ExtractedIDMatrix:
type fieldMapper = func([]uint64) (interface{}, error)
fields := make([]ExtractedTableField, len(result.Fields))
mappers := make([]fieldMapper, len(result.Fields))
for i, v := range result.Fields {
field := idx.Field(v)
if field == nil {
return nil, ErrFieldNotFound
}
typ := field.Type()
fields[i] = ExtractedTableField{
Name: v,
Type: typ,
}
var mapper fieldMapper
switch typ {
case FieldTypeBool:
mapper = func(ids []uint64) (interface{}, error) {
switch len(ids) {
case 0:
return nil, nil
case 1:
switch ids[0] {
case 0:
return false, nil
case 1:
return true, nil
default:
return nil, errors.Errorf("invalid ID for boolean %q: %d", field.Name(), ids[0])
}
default:
return nil, errors.Errorf("boolean %q has too many values: %v", field.Name(), ids)
}
}
case FieldTypeSet, FieldTypeTime:
if field.Keys() {
translator := field.TranslateStore()
mapper = func(ids []uint64) (interface{}, error) {
return translator.TranslateIDs(ids)
}
} else {
mapper = func(ids []uint64) (interface{}, error) {
if ids == nil {
ids = []uint64{}
}
return ids, nil
}
}
case FieldTypeMutex:
if field.Keys() {
translator := field.TranslateStore()
mapper = func(ids []uint64) (interface{}, error) {
switch len(ids) {
case 0:
return nil, nil
case 1:
return translator.TranslateID(ids[0])
default:
return nil, errors.Errorf("mutex %q has too many values: %v", field.Name(), ids)
}
}
} else {
mapper = func(ids []uint64) (interface{}, error) {
switch len(ids) {
case 0:
return nil, nil
case 1:
return ids[0], nil
default:
return nil, errors.Errorf("mutex %q has too many values: %v", field.Name(), ids)
}
}
}
case FieldTypeInt:
mapper = func(ids []uint64) (interface{}, error) {
switch len(ids) {
case 0:
return nil, nil
case 1:
return int64(ids[0]), nil
default:
return nil, errors.Errorf("BSI field %q has too many values: %v", field.Name(), ids)
}
}
case FieldTypeDecimal:
scale := field.Options().Scale
mapper = func(ids []uint64) (interface{}, error) {
switch len(ids) {
case 0:
return nil, nil
case 1:
return pql.NewDecimal(int64(ids[0]), scale), nil
default:
return nil, errors.Errorf("BSI field %q has too many values: %v", field.Name(), ids)
}
}
default:
return nil, errors.Errorf("field type %q not yet supported", typ)
}
mappers[i] = mapper
}
var translateCol func(uint64) (KeyOrID, error)
if idx.keys {
translateCol = func(id uint64) (KeyOrID, error) {
return KeyOrID{Keyed: true, Key: idSet[id]}, nil
}
} else {
translateCol = func(id uint64) (KeyOrID, error) {
return KeyOrID{ID: id}, nil
}
}
cols := make([]ExtractedTableColumn, len(result.Columns))
colData := make([]interface{}, len(cols)*len(result.Fields))
for i, col := range result.Columns {
data := colData[i*len(result.Fields) : (i+1)*len(result.Fields) : (i+1)*len(result.Fields)]
for j, rows := range col.Rows {
v, err := mappers[j](rows)
if err != nil {
return nil, errors.Wrap(err, "translating extracted table value")
}
data[j] = v
}
colTrans, err := translateCol(col.ColumnID)
if err != nil {
return nil, errors.Wrap(err, "translating column ID in extracted table")
}
cols[i] = ExtractedTableColumn{
Column: colTrans,
Rows: data,
}
}
return ExtractedTable{
Fields: fields,
Columns: cols,
}, nil
}
return result, nil

View file

@ -61,6 +61,25 @@ func getTempDirString() (td *string) {
return td
}
func TestExecutor_Execute_ConstRow(t *testing.T) {
c := test.MustRunCluster(t, 2)
defer c.Close()
c.CreateField(t, "i", pilosa.IndexOptions{}, "h")
c.ImportBits(t, "i", "h", [][2]uint64{
{1, 2},
{3, 4},
{5, 6},
})
resp := c.Query(t, "i", `ConstRow(columns=[2,6])`)
expect := []uint64{2, 6}
got := resp.Results[0].(*pilosa.Row).Columns()
if !reflect.DeepEqual(expect, got) {
t.Errorf("expected %v but got %v", expect, got)
}
}
// Ensure a row query can be executed.
func TestExecutor_Execute_Row(t *testing.T) {
t.Run("RowIDColumnID", func(t *testing.T) {
@ -3631,10 +3650,96 @@ func TestExecutor_Execute_FieldValue(t *testing.T) {
}
}
// Ensure a Limit query can be executed.
func TestExecutor_Execute_Limit(t *testing.T) {
c := test.MustRunCluster(t, 1)
defer c.Close()
c.CreateField(t, "i", pilosa.IndexOptions{TrackExistence: true}, "f")
c.ImportBits(t, "i", "f", [][2]uint64{
{1, 0},
{1, 1},
{1, ShardWidth + 1},
})
columns := []uint64{0, 1, ShardWidth + 1}
// Test with only a limit specified.
t.Run("Limit", func(t *testing.T) {
for limit := 0; limit < 5; limit++ {
expect := columns
if limit < len(expect) {
expect = expect[:limit]
}
resp := c.Query(t, "i", fmt.Sprintf("Limit(All(), limit=%d)", limit))
if len(resp.Results) != 1 {
t.Fatalf("limit=%d: expected 1 result but got %v", limit, resp.Results)
}
row, ok := resp.Results[0].(*pilosa.Row)
if !ok {
t.Fatalf("limit=%d: expected a row result but got %T", limit, resp.Results[0])
}
got := row.Columns()
if !reflect.DeepEqual(expect, got) {
t.Errorf("limit=%d: expected %v but got %v", limit, expect, got)
}
}
})
// Test with only an offset specified.
t.Run("Offset", func(t *testing.T) {
for offset := 0; offset < 5; offset++ {
expect := []uint64{}
if offset <= len(columns) {
expect = columns[offset:]
}
resp := c.Query(t, "i", fmt.Sprintf("Limit(All(), offset=%d)", offset))
if len(resp.Results) != 1 {
t.Fatalf("offset=%d: expected 1 result but got %v", offset, resp.Results)
}
row, ok := resp.Results[0].(*pilosa.Row)
if !ok {
t.Fatalf("offset=%d: expected a row result but got %T", offset, resp.Results[0])
}
got := row.Columns()
if !reflect.DeepEqual(expect, got) {
t.Errorf("offset=%d: expected %v but got %v", offset, expect, got)
}
}
})
// Test with a limit and offset specified.
t.Run("LimitOffset", func(t *testing.T) {
for limit := 0; limit < 5; limit++ {
for offset := 0; offset < 5; offset++ {
expect := []uint64{}
if offset <= len(columns) {
expect = columns[offset:]
}
if limit < len(expect) {
expect = expect[:limit]
}
resp := c.Query(t, "i", fmt.Sprintf("Limit(All(), limit=%d, offset=%d)", limit, offset))
if len(resp.Results) != 1 {
t.Fatalf("limit=%d,offset=%d: expected 1 result but got %v", limit, offset, resp.Results)
}
row, ok := resp.Results[0].(*pilosa.Row)
if !ok {
t.Fatalf("limit=%d,offset=%d: expected a row result but got %T", limit, offset, resp.Results[0])
}
got := row.Columns()
if !reflect.DeepEqual(expect, got) {
t.Errorf("limit=%d,offset=%d: expected %v but got %v", limit, offset, expect, got)
}
}
}
})
}
// Ensure an all query can be executed.
func TestExecutor_Execute_All(t *testing.T) {
skipForRBF(t)
t.Run("ColumnID", func(t *testing.T) {
c := test.MustRunCluster(t, 1)
defer c.Close()
@ -4235,6 +4340,291 @@ func benchmarkExistence(nn bool, b *testing.B) {
func BenchmarkExecutor_Existence_True(b *testing.B) { benchmarkExistence(true, b) }
func BenchmarkExecutor_Existence_False(b *testing.B) { benchmarkExistence(false, b) }
func TestExecutor_Execute_Extract(t *testing.T) {
c := test.MustRunCluster(t, 3)
defer c.Close()
c.CreateField(t, "i", pilosa.IndexOptions{TrackExistence: true}, "set")
c.ImportBits(t, "i", "set", [][2]uint64{
{0, 1},
{0, 2},
{3, 1},
{4, 1},
{4, 4 * ShardWidth},
{5, ShardWidth},
})
c.Query(t, "i", fmt.Sprintf("Clear(%d, set=5)", ShardWidth))
c.CreateField(t, "i", pilosa.IndexOptions{TrackExistence: true}, "keyset", pilosa.OptFieldKeys())
c.Query(t, "i", `
Set(0, keyset="h")
Set(1, keyset="xyzzy")
Set(0, keyset="plugh")
`)
c.CreateField(t, "i", pilosa.IndexOptions{TrackExistence: true}, "mutex", pilosa.OptFieldTypeMutex(pilosa.CacheTypeRanked, 5000))
c.ImportBits(t, "i", "mutex", [][2]uint64{
{0, 1},
{0, 2},
{4, 4 * ShardWidth},
})
c.CreateField(t, "i", pilosa.IndexOptions{TrackExistence: true}, "keymutex", pilosa.OptFieldKeys(), pilosa.OptFieldTypeMutex(pilosa.CacheTypeRanked, 5000))
c.Query(t, "i", `
Set(0, keymutex="h")
Set(1, keymutex="xyzzy")
Set(3, keymutex="plugh")
`)
c.CreateField(t, "i", pilosa.IndexOptions{TrackExistence: true}, "time", pilosa.OptFieldTypeTime("YMDH"))
c.Query(t, "i", `
Set(0, time=1, 2016-01-01T00:00)
Set(1, time=2, 2017-01-01T00:00)
Set(3, time=3, 2018-01-01T00:00)
`)
c.CreateField(t, "i", pilosa.IndexOptions{TrackExistence: true}, "keytime", pilosa.OptFieldKeys(), pilosa.OptFieldTypeTime("YMDH"))
c.Query(t, "i", `
Set(0, keytime="h", 2016-01-01T00:00)
Set(1, keytime="xyzzy", 2017-01-01T00:00)
Set(0, keytime="plugh", 2018-01-01T00:00)
`)
c.CreateField(t, "i", pilosa.IndexOptions{TrackExistence: true}, "bsint", pilosa.OptFieldTypeInt(-100, 100))
c.Query(t, "i", `
Set(0, bsint=1)
Set(1, bsint=-1)
Set(3, bsint=2)
`)
c.CreateField(t, "i", pilosa.IndexOptions{TrackExistence: true}, "bsidecimal", pilosa.OptFieldTypeDecimal(2))
c.Query(t, "i", `
Set(0, bsidecimal=0.01)
Set(1, bsidecimal=1.00)
Set(3, bsidecimal=-1.01)
`)
c.CreateField(t, "i", pilosa.IndexOptions{TrackExistence: true}, "bool", pilosa.OptFieldTypeBool())
c.Query(t, "i", `
Set(0, bool=true)
Set(1, bool=false)
Set(3, bool=true)
`)
resp := c.Query(t, "i", `Extract(All(), Rows(set), Rows(keyset), Rows(mutex), Rows(keymutex), Rows(time), Rows(keytime), Rows(bsint), Rows(bsidecimal), Rows(bool))`)
expect := []interface{}{
pilosa.ExtractedTable{
Fields: []pilosa.ExtractedTableField{
{
Name: "set",
Type: pilosa.FieldTypeSet,
},
{
Name: "keyset",
Type: pilosa.FieldTypeSet,
},
{
Name: "mutex",
Type: pilosa.FieldTypeMutex,
},
{
Name: "keymutex",
Type: pilosa.FieldTypeMutex,
},
{
Name: "time",
Type: pilosa.FieldTypeTime,
},
{
Name: "keytime",
Type: pilosa.FieldTypeTime,
},
{
Name: "bsint",
Type: pilosa.FieldTypeInt,
},
{
Name: "bsidecimal",
Type: pilosa.FieldTypeDecimal,
},
{
Name: "bool",
Type: pilosa.FieldTypeBool,
},
},
Columns: []pilosa.ExtractedTableColumn{
{
Column: pilosa.KeyOrID{ID: 0},
Rows: []interface{}{
[]uint64{},
[]string{
"h",
"plugh",
},
nil,
"h",
[]uint64{
1,
},
[]string{
"h",
"plugh",
},
int64(1),
pql.NewDecimal(1, 2),
true,
},
},
{
Column: pilosa.KeyOrID{ID: 1},
Rows: []interface{}{
[]uint64{
0,
3,
4,
},
[]string{
"xyzzy",
},
uint64(0),
"xyzzy",
[]uint64{
2,
},
[]string{
"xyzzy",
},
int64(-1),
pql.NewDecimal(100, 2),
false,
},
},
{
Column: pilosa.KeyOrID{ID: 2},
Rows: []interface{}{
[]uint64{
0,
},
[]string{},
uint64(0),
nil,
[]uint64{},
[]string{},
nil,
nil,
nil,
},
},
{
Column: pilosa.KeyOrID{ID: 3},
Rows: []interface{}{
[]uint64{},
[]string{},
nil,
"plugh",
[]uint64{
3,
},
[]string{},
int64(2),
pql.NewDecimal(-101, 2),
true,
},
},
{
Column: pilosa.KeyOrID{ID: ShardWidth},
Rows: []interface{}{
[]uint64{},
[]string{},
nil,
nil,
[]uint64{},
[]string{},
nil,
nil,
nil,
},
},
{
Column: pilosa.KeyOrID{ID: 4 * ShardWidth},
Rows: []interface{}{
[]uint64{
4,
},
[]string{},
uint64(4),
nil,
[]uint64{},
[]string{},
nil,
nil,
nil,
},
},
},
},
}
if !reflect.DeepEqual(expect, resp.Results) {
t.Errorf("expected %v but got %v", expect, resp.Results)
}
}
func TestExecutor_Execute_Extract_Keyed(t *testing.T) {
c := test.MustRunCluster(t, 3)
defer c.Close()
c.CreateField(t, "i", pilosa.IndexOptions{TrackExistence: true, Keys: true}, "set")
c.Query(t, "i", `
Set("h", set=1)
Set("h", set=2)
Set("xyzzy", set=2)
Set("plugh", set=1)
Clear("plugh", set=1)
`)
resp := c.Query(t, "i", `Extract(All(), Rows(set))`)
expect := []interface{}{
pilosa.ExtractedTable{
Fields: []pilosa.ExtractedTableField{
{
Name: "set",
Type: "set",
},
},
Columns: []pilosa.ExtractedTableColumn{
{
Column: pilosa.KeyOrID{Keyed: true, Key: "plugh"},
Rows: []interface{}{
[]uint64{},
},
},
{
Column: pilosa.KeyOrID{Keyed: true, Key: "h"},
Rows: []interface{}{
[]uint64{
1,
2,
},
},
},
{
Column: pilosa.KeyOrID{Keyed: true, Key: "xyzzy"},
Rows: []interface{}{
[]uint64{
2,
},
},
},
},
},
}
if !reflect.DeepEqual(expect, resp.Results) {
t.Errorf("expected %v but got %v", expect, resp.Results)
}
}
func TestExecutor_Execute_Rows(t *testing.T) {
c := test.MustRunCluster(t, 3)
defer c.Close()
@ -4395,8 +4785,6 @@ func TestExecutor_Execute_Query_Error(t *testing.T) {
}
func TestExecutor_GroupByStrings(t *testing.T) {
skipForRBF(t)
c := test.MustRunCluster(t, 1)
defer c.Close()
c.CreateField(t, "istring", pilosa.IndexOptions{Keys: true}, "generals", pilosa.OptFieldKeys())
@ -4853,8 +5241,6 @@ func sameStringSlice(x, y []string) bool {
}
func TestExecutor_Execute_GroupBy(t *testing.T) {
skipForRBF(t)
groupByTest := func(t *testing.T, clusterSize int) {
c := test.MustRunCluster(t, 1)
defer c.Close()

View file

@ -748,6 +748,7 @@ fileLoop:
}()
name := filepath.Base(fi.Name())
f.holder.Logger.Debugf("open index/field/view: %s/%s/%s", f.index, f.name, fi.Name())
view := f.newView(f.viewPath(name), name)
if err := view.open(); err != nil {
return fmt.Errorf("opening view: view=%s, err=%s", view.name, err)

View file

@ -258,7 +258,7 @@ func (f *fragment) Open() error {
f.checksums = make(map[int][]byte)
// Read last bit to determine max row.
tx := f.idx.Txf.NewTx(Txo{Write: !writable, Index: f.idx, Fragment: f})
tx := f.idx.Txf.NewTx(Txo{Write: false, Index: f.idx, Fragment: f})
defer tx.Rollback()
return f.calculateMaxRowID(tx)
}(); err != nil {

View file

@ -1604,10 +1604,8 @@ func TestFragment_TopN_CacheSize(t *testing.T) {
defer f.Clean(t)
// Obtain transaction.
tx := index.Txf.NewTx(Txo{Write: writable, Index: index, Fragment: f})
f.txTestingOnly = tx
defer tx.Rollback() // okay to call 2x. f.Clean(t) will do Rollback() too.
tx := index.Txf.NewTx(Txo{Write: writable, Index: index})
defer tx.Rollback()
// Set bits on various rows.
f.mustSetBits(tx, 100, 1, 2, 3)
@ -1812,7 +1810,7 @@ func TestFragment_RankCache_Persistence(t *testing.T) {
}
// Obtain transaction.
tx := index.Txf.NewTx(Txo{Write: writable, Index: index, Fragment: f})
tx := index.Txf.NewTx(Txo{Write: writable, Index: index})
defer tx.Rollback()
// Set bits on the fragment.
@ -5324,8 +5322,6 @@ func check(t *testing.T, tx Tx, f *fragment, exp map[uint64]map[uint64]struct{})
}
func TestImportValueConcurrent(t *testing.T) {
skipForRBF(t)
f, idx := mustOpenBSIFragment("i", "f", viewBSIGroupPrefix+"foo", 0)
switch idx.Txf.TxType() {
case blueGreenBadgerRoaring, blueGreenRoaringBadger:

File diff suppressed because it is too large Load diff

View file

@ -19,6 +19,53 @@ message RowIdentifiers {
repeated string Keys = 2;
}
message IDList {
repeated uint64 IDs = 1;
}
message ExtractedIDColumn {
uint64 ID = 1;
repeated IDList Vals = 2;
}
message ExtractedIDMatrix {
repeated string Fields = 1;
repeated ExtractedIDColumn Columns = 2;
}
message KeyList {
repeated string Keys = 1;
}
message ExtractedTableValue {
oneof Value {
IDList IDs = 1;
KeyList Keys = 2;
int64 BSIValue = 3;
uint64 MutexID = 4;
string MutexKey = 5;
bool Bool = 6;
}
}
message ExtractedTableColumn {
oneof KeyOrID {
string Key = 1;
uint64 ID = 2;
}
repeated ExtractedTableValue Values = 3;
}
message ExtractedTableField {
string Name = 1;
string Type = 2;
}
message ExtractedTable {
repeated ExtractedTableField Fields = 1;
repeated ExtractedTableColumn Columns = 2;
}
message Pair {
uint64 ID = 1;
string Key = 3;
@ -112,6 +159,8 @@ message QueryResult {
SignedRow SignedRow = 10;
PairsField PairsField = 11;
PairField PairField = 12;
ExtractedIDMatrix ExtractedIDMatrix = 13;
ExtractedTable ExtractedTable = 14;
}
message ImportRequest {

View file

@ -15,7 +15,6 @@
package pilosa_test
import (
"os"
"strings"
"testing"
@ -56,9 +55,3 @@ func TestAddressWithDefaults(t *testing.T) {
}
}
}
func skipForRBF(tb testing.TB) {
if os.Getenv("PILOSA_TXSRC") == "rbf" {
tb.Skip("skip for RBF")
}
}

View file

@ -393,7 +393,22 @@ var callInfoByFunc = map[string]callInfo{
},
"Union": {allowUnknown: false},
"UnionRows": {allowUnknown: false},
"Xor": {allowUnknown: false},
"Extract": {allowUnknown: false},
"Limit": {
allowUnknown: false,
prototypes: map[string]interface{}{
"limit": int64(0),
"offset": int64(0),
},
},
"Xor": {allowUnknown: false},
"ConstRow": {
allowUnknown: false,
prototypes: map[string]interface{}{
"columns": []interface{}{},
},
},
// things that take _field
"TopN": allowUnderField,

View file

@ -187,6 +187,7 @@ func (db *DB) checkpoint() error {
// Loop over each transaction
walID++
pageMap := immutable.NewMap(&uint32Hasher{})
var maxCheckpointedWALID int64
for {
// Determine last page of transaction.
metaWALID, metaFlags, err := db.findNextWALMetaPage(walID)
@ -240,20 +241,27 @@ func (db *DB) checkpoint() error {
if err := db.writePage(pgno, page); err != nil {
return err
}
// Track highest WALID that has been checkpointed back to disk.
if IsMetaPage(page) {
maxCheckpointedWALID = walID
}
}
}
// Remove WAL segments that have been checkpointed.
for len(db.segments) > 1 {
segment := db.segments[0]
if minActiveWALID != 0 && segment.MaxWALID() >= minActiveWALID {
break
}
if maxCheckpointedWALID != 0 {
for len(db.segments) > 1 {
segment := db.segments[0]
if segment.MaxWALID() >= maxCheckpointedWALID {
break
}
if err := segment.Close(); err != nil {
return err
if err := segment.Close(); err != nil {
return err
}
db.segments, db.segments[0] = db.segments[1:], nil
}
db.segments, db.segments[0] = db.segments[1:], nil
}
db.pageMap = pageMap
@ -424,13 +432,13 @@ func (db *DB) addWALSegment() error {
func (db *DB) Close() (err error) {
// TODO(bbj): Add wait group to hang until last Tx is complete.
db.mu.Lock()
defer db.mu.Unlock()
// Wait for writer lock.
db.rwmu.Lock()
defer db.rwmu.Unlock()
db.mu.Lock()
defer db.mu.Unlock()
db.opened = false
// Close mmap handle.
@ -546,6 +554,11 @@ func (db *DB) initFreelistPage() error {
func (db *DB) Begin(writable bool) (_ *Tx, err error) {
// TODO(BBJ): Acquire write lock if writable.
// Ensure only one writable transaction at a time.
if writable {
db.rwmu.Lock()
}
db.mu.Lock()
defer db.mu.Unlock()
@ -555,11 +568,6 @@ func (db *DB) Begin(writable bool) (_ *Tx, err error) {
tx := &Tx{db: db, pageMap: db.pageMap, writable: writable}
// Ensure only one writable transaction at a time.
if tx.writable {
db.rwmu.Lock()
}
// Copy meta page into transaction's buffer.
// This page is only written at the end of a dirty transaction.
page, err := db.readPage(db.pageMap, 0)
@ -590,8 +598,10 @@ func (db *DB) removeTx(tx *Tx) error {
// Write pages from WAL to DB.
// TODO(bbj): Move this to an async goroutine.
if err := db.checkpoint(); err != nil {
return err
if tx.writable {
if err := db.checkpoint(); err != nil {
return err
}
}
delete(tx.db.txs, tx)

View file

@ -106,18 +106,8 @@ func (tx *Tx) Rollback() {
}
}
// turn on these error checks! we see
// panic: cannot find segment containing WAL page: 1
// when running go test -v
// TestCursor_FirstNext_Quick/6
//
//panicOn(tx.db.checkpoint())
//panicOn(tx.db.removeTx(tx))
_ = tx.db.checkpoint()
// Disconnect transaction from DB.
_ = tx.db.removeTx(tx)
panicOn(tx.db.removeTx(tx))
}
// Root returns the root page number for a bitmap. Returns 0 if the bitmap does not exist.

View file

@ -720,8 +720,7 @@ func (s *Server) receiveMessage(m Message) error {
if f == nil {
return fmt.Errorf("local field not found: %s", obj.Field)
}
_, _, err := f.createViewIfNotExistsBase(obj.View)
if err != nil {
if _, _, err := f.createViewIfNotExistsBase(obj.View); err != nil {
return err
}
case *DeleteViewMessage:

View file

@ -166,8 +166,7 @@ var workQueue = make(chan struct{}, runtime.NumCPU()*2)
// replaces v.openFragments() with Tx generic code.
func (v *view) openFragmentsInTx() error {
tx := v.idx.Txf.NewTx(Txo{Write: !writable, Index: v.idx})
tx := v.idx.Txf.NewTx(Txo{Write: false, Index: v.idx})
defer tx.Rollback()
shards, err := tx.SliceOfShards(v.index, v.field, v.name, v.path)