From f2c096040a20593897d6bec92df2d3dadd7a8f22 Mon Sep 17 00:00:00 2001 From: Yuce Tekol Date: Mon, 22 Oct 2018 15:47:29 +0300 Subject: [PATCH 01/11] updated api docs --- docs/api-reference.md | 124 +++++++++++++++++++++++++++++++++++++++--- 1 file changed, 116 insertions(+), 8 deletions(-) diff --git a/docs/api-reference.md b/docs/api-reference.md index 216a01bca..e59d06822 100644 --- a/docs/api-reference.md +++ b/docs/api-reference.md @@ -17,9 +17,42 @@ Returns the schema of all indexes in JSON. curl -XGET localhost:10101/index ``` ``` response -{"indexes":[{"name":"user","fields":[{"name":"event","options":{"type":"time","timeQuantum":"YMD","keys":false}}]}]} +{ + "indexes": [ + { + "fields": [ + { + "name": "event", + "options": { + "keys": false, + "timeQuantum": "YMD", + "type": "time" + } + }, + { + "name": "language", + "options": { + "cacheSize": 50000, + "cacheType": "ranked", + "keys": false, + "type": "set" + } + } + ], + "name": "user", + "options": { + "keys": false, + "trackExistence": true + } + } + ] +} ``` +`GET /schema` + +Is equivalent to `GET /index` and returns the same response. + ### List index schema `GET /index/` @@ -30,7 +63,23 @@ Returns the schema of the specified index in JSON. curl -XGET localhost:10101/index/user ``` ``` response -{"name":"user","fields":[{"name":"event","options":{"type":"time","timeQuantum":"YMD","keys":false}}]} +{ + "fields": [ + { + "name": "event", + "options": { + "keys": false, + "timeQuantum": "YMD", + "type": "time" + } + } + ], + "name": "user", + "options": { + "keys": false, + "trackExistence": true + } +} ``` ### Create index @@ -39,8 +88,13 @@ curl -XGET localhost:10101/index/user Creates an index with the given name. +The request payload is in JSON, and may contain the `options` field. The `options` field is a JSON object with the following options: + +* `keys` (bool): Enables using column keys instead of column IDs. +* `trackExistence` (bool): Enables or disables existence tracking on the index. Required for [Not](../query-language/#not) queries. It is `true` by default. Note that disabling track existence improves query performance. + ``` request -curl -XPOST localhost:10101/index/user +curl -XPOST localhost:10101/index/user -d '{"options":{"keys":true}}' ``` ``` response {"success":true} @@ -71,7 +125,16 @@ curl localhost:10101/index/user/query \ -d 'Row(language=5)' ``` ``` response -{"results":[{"attrs":{},"columns":[100]}]} +{ + "results": [ + { + "attrs": {}, + "columns": [ + 100 + ] + } + ] +} ``` In order to send protobuf binaries in the request and response, set `Content-Type` and `Accept` headers to: `application/x-protobuf`. @@ -87,8 +150,22 @@ curl "localhost:10101/index/user/query?columnAttrs=true&shards=0,1" \ ``` ``` response { - "results":[{"attrs":{},"columns":[100]}], - "columnAttrs":[{"id":100,"attrs":{"name":"Klingon"}}] + "columnAttrs": [ + { + "attrs": { + "name": "Klingon" + }, + "id": 100 + } + ], + "results": [ + { + "attrs": {}, + "columns": [ + 100 + ] + } + ] } ``` @@ -100,7 +177,12 @@ By default, all bits and attributes (*for `Row` queries only*) are returned. In Creates a field in the given index with the given name. -The request payload is in JSON, and may contain the `options` field. The `options` field is a JSON object which must contain a `type` along with the corresponding configuration options. +The request payload is in JSON, and may contain the `options` field. The `options` field is a JSON object which must contain a `type` and optionally the `keys`: + +* `keys` (bool): Enables using column keys instead of column IDs. +* `type` (string): Sets the field type and type options. + +Valid `type`s and correspondonding options are listed below: * `set` * `cacheType` (string): [ranked](../data-model/#ranked) or [LRU](../data-model/#lru) caching on this field. Default is `ranked`. @@ -171,6 +253,33 @@ curl -XGET localhost:10101/version {"version":"v0.6.0"} ``` +### Get status + +`GET /status` + +Returns the status of nodes in a Pilosa cluster. + +```request +curl -XGET localhost:10101/status +``` +```response +{ + "localID": "d3369125-29d8-4305-a351-b4474d14a542", + "nodes": [ + { + "id": "d3369125-29d8-4305-a351-b4474d14a542", + "isCoordinator": true, + "uri": { + "host": "localhost", + "port": 10101, + "scheme": "http" + } + } + ], + "state": "NORMAL" +} +``` + ### Recalculate Caches `POST /recalculate-caches` @@ -187,4 +296,3 @@ curl -XPOST localhost:10101/recalculate-caches ``` Response: `204 No Content` - From 5060e6ee4745e8947db2904a562596358004c562 Mon Sep 17 00:00:00 2001 From: Yuce Tekol Date: Mon, 22 Oct 2018 17:24:21 +0300 Subject: [PATCH 02/11] updated --- docs/api-reference.md | 92 ++++++++++++++++++++++--------------------- 1 file changed, 47 insertions(+), 45 deletions(-) diff --git a/docs/api-reference.md b/docs/api-reference.md index e59d06822..b8e4d344b 100644 --- a/docs/api-reference.md +++ b/docs/api-reference.md @@ -11,47 +11,7 @@ nav = [] `GET /index` -Returns the schema of all indexes in JSON. - -``` request -curl -XGET localhost:10101/index -``` -``` response -{ - "indexes": [ - { - "fields": [ - { - "name": "event", - "options": { - "keys": false, - "timeQuantum": "YMD", - "type": "time" - } - }, - { - "name": "language", - "options": { - "cacheSize": 50000, - "cacheType": "ranked", - "keys": false, - "type": "set" - } - } - ], - "name": "user", - "options": { - "keys": false, - "trackExistence": true - } - } - ] -} -``` - -`GET /schema` - -Is equivalent to `GET /index` and returns the same response. +Is equivalent to `GET /schema` and returns the same response. ### List index schema @@ -91,7 +51,7 @@ Creates an index with the given name. The request payload is in JSON, and may contain the `options` field. The `options` field is a JSON object with the following options: * `keys` (bool): Enables using column keys instead of column IDs. -* `trackExistence` (bool): Enables or disables existence tracking on the index. Required for [Not](../query-language/#not) queries. It is `true` by default. Note that disabling track existence improves query performance. +* `trackExistence` (bool): Enables or disables existence tracking on the index. Required for [Not](../query-language/#not) queries. It is `true` by default. ``` request curl -XPOST localhost:10101/index/user -d '{"options":{"keys":true}}' @@ -177,10 +137,10 @@ By default, all bits and attributes (*for `Row` queries only*) are returned. In Creates a field in the given index with the given name. -The request payload is in JSON, and may contain the `options` field. The `options` field is a JSON object which must contain a `type` and optionally the `keys`: +The request payload is in JSON, and may contain the `options` field. The `options` field is a JSON object which must contain a `type`: -* `keys` (bool): Enables using column keys instead of column IDs. * `type` (string): Sets the field type and type options. +* `keys` (bool): Enables using column keys instead of column IDs (optional). Valid `type`s and correspondonding options are listed below: @@ -240,6 +200,48 @@ curl -XDELETE localhost:10101/index/user/field/language {"success":true} ``` +### List all index schemas + +`GET /schema` + +Returns the schema of all indexes in JSON. + +``` request +curl -XGET localhost:10101/index +``` +``` response +{ + "indexes": [ + { + "fields": [ + { + "name": "event", + "options": { + "keys": false, + "timeQuantum": "YMD", + "type": "time" + } + }, + { + "name": "language", + "options": { + "cacheSize": 50000, + "cacheType": "ranked", + "keys": false, + "type": "set" + } + } + ], + "name": "user", + "options": { + "keys": false, + "trackExistence": true + } + } + ] +} +``` + ### Get version `GET /version` @@ -257,7 +259,7 @@ curl -XGET localhost:10101/version `GET /status` -Returns the status of nodes in a Pilosa cluster. +Returns the status of the cluster. ```request curl -XGET localhost:10101/status From 4eed160b2de55f0377135bb3405e933676a0f17a Mon Sep 17 00:00:00 2001 From: Matthew Jaffee Date: Mon, 22 Oct 2018 23:30:13 -0500 Subject: [PATCH 03/11] wip on adding import docs --- docs/api-reference.md | 36 ++++++++++++++++++++++++++++++++++++ docs/getting-started.md | 5 +++++ docs/query-language.md | 4 ++++ 3 files changed, 45 insertions(+) diff --git a/docs/api-reference.md b/docs/api-reference.md index 216a01bca..4a49e3515 100644 --- a/docs/api-reference.md +++ b/docs/api-reference.md @@ -94,6 +94,40 @@ curl "localhost:10101/index/user/query?columnAttrs=true&shards=0,1" \ By default, all bits and attributes (*for `Row` queries only*) are returned. In order to suppress returning bits, set `excludeBits` query argument to `true`; to suppress returning attributes, set `excludeAttrs` query argument to `true`. +### Import Data + +`POST /index//field/` + +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 +is not usually necessary to use this endpoint directly. See the documentation for +imports for +Go, +Java, +and Python. + +The request payload is protobuf encoded with the following schema. The row or +column Keys fields are used if the field or index is configured for keys +respectively. Otherwise, the RowIDs and ColumnIDs fields are used. They must +have the same number of items, and each index into those two lists represents a +particular bit to be set. Timestamps are optional, but if they exist must also +contain the same number of items as rows and columns. The column IDs must all be +in the shard specified in the request. + +``` +message ImportRequest { + string Index = 1; + string Field = 2; + uint64 Shard = 3; + repeated uint64 RowIDs = 4; + repeated uint64 ColumnIDs = 5; + repeated string RowKeys = 7; + repeated string ColumnKeys = 8; + repeated int64 Timestamps = 6; +} +``` + + ### Create field `POST /index//field/` @@ -158,6 +192,8 @@ curl -XDELETE localhost:10101/index/user/field/language {"success":true} ``` + + ### Get version `GET /version` diff --git a/docs/getting-started.md b/docs/getting-started.md index 563176c62..257e77628 100644 --- a/docs/getting-started.md +++ b/docs/getting-started.md @@ -89,6 +89,11 @@ The `language` is a `set` field, but since the default field type is `set`, we d #### Import Data From CSV Files +
+

For demonstration purposes, we're using Pilosa's built in utility to import specially formatted CSV files. For more general usage, see how the various client libraries expose the bulk import functionality in Go, Java, and Python.

+
+ + Download the `stargazer.csv` and `language.csv` files here: ``` diff --git a/docs/query-language.md b/docs/query-language.md index 9bc231e6f..0f8feb795 100644 --- a/docs/query-language.md +++ b/docs/query-language.md @@ -67,6 +67,10 @@ Set(, =, [TIMESTAMP]) `Set` assigns a value of 1 to a bit in the binary matrix, thus associating the given row (the `` value) in the given field with the given column. +
+

While using "Set" in PQL is a convenient way to get familiar with Pilosa, it's almost always better to use the import functionality in the Go, Java, and Python clients to ingest lots of data.

+
+ **Result Type:** boolean A return value of `true` indicates that the bit was changed to 1. From 376c2d61ffa7c90df3b90d0251f892dd11008a22 Mon Sep 17 00:00:00 2001 From: Matt Jaffee Date: Tue, 23 Oct 2018 13:50:35 -0500 Subject: [PATCH 04/11] address CR feedback --- docs/api-reference.md | 16 ++++++++-------- 1 file changed, 8 insertions(+), 8 deletions(-) diff --git a/docs/api-reference.md b/docs/api-reference.md index 4a49e3515..f5f06fd57 100644 --- a/docs/api-reference.md +++ b/docs/api-reference.md @@ -96,7 +96,7 @@ By default, all bits and attributes (*for `Row` queries only*) are returned. In ### Import Data -`POST /index//field/` +`POST /index//field//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 @@ -106,13 +106,13 @@ imports for Java, and Python. -The request payload is protobuf encoded with the following schema. The row or -column Keys fields are used if the field or index is configured for keys -respectively. Otherwise, the RowIDs and ColumnIDs fields are used. They must -have the same number of items, and each index into those two lists represents a -particular bit to be set. Timestamps are optional, but if they exist must also -contain the same number of items as rows and columns. The column IDs must all be -in the shard specified in the request. +The request payload is protobuf encoded with the following schema. The RowKeys +and/or ColumnKeys fields are used if the pilosa field or index are configured +for keys respectively. Otherwise, the RowIDs and ColumnIDs fields are used. They +must have the same number of items, and each index into those two lists +represents a particular bit to be set. Timestamps are optional, but if they +exist must also contain the same number of items as rows and columns. The +column IDs must all be in the shard specified in the request. ``` message ImportRequest { From caf8e067129b8c7fdcbd44c6b1c937dcc67cbc74 Mon Sep 17 00:00:00 2001 From: Travis Turner Date: Mon, 22 Oct 2018 13:35:32 -0500 Subject: [PATCH 05/11] add `clear` functional option for imports --- api.go | 36 +++++- client.go | 8 +- cmd/import.go | 1 + cmd/import_test.go | 11 ++ ctl/import.go | 7 +- ctl/import_test.go | 67 +++++++--- field.go | 22 +++- fragment.go | 26 ++-- fragment_internal_test.go | 250 +++++++++++++++++++++++++++++++++++--- http/client.go | 46 +++++-- http/handler.go | 10 +- 11 files changed, 413 insertions(+), 71 deletions(-) diff --git a/api.go b/api.go index 0b27d19a8..bb900db4e 100644 --- a/api.go +++ b/api.go @@ -670,12 +670,36 @@ func (api *API) FieldAttrDiff(_ context.Context, indexName string, fieldName str return attrs, nil } +// ImportOptions holds the options for the API.Import method. +type ImportOptions struct { + Clear bool +} + +// ImportOption is a functional option type for API.Import +type ImportOption func(*ImportOptions) error + +func OptImportOptionsClear(c bool) ImportOption { + return func(o *ImportOptions) error { + o.Clear = c + return nil + } +} + // Import bulk imports data into a particular index,field,shard. -func (api *API) Import(_ context.Context, req *ImportRequest) error { +func (api *API) Import(_ context.Context, req *ImportRequest, opts ...ImportOption) error { if err := api.validate(apiImport); err != nil { return errors.Wrap(err, "validating api method") } + // Set up import options. + options := &ImportOptions{} + for _, opt := range opts { + err := opt(options) + if err != nil { + return errors.Wrap(err, "applying option") + } + } + index := api.holder.Index(req.Index) if index == nil { return newNotFoundError(ErrIndexNotFound) @@ -717,13 +741,15 @@ func (api *API) Import(_ context.Context, req *ImportRequest) error { } // Import columnIDs into existence field. - if err := importExistenceColumns(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 errors.Wrap(err, "importing existence columns") + if !options.Clear { + if err := importExistenceColumns(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 errors.Wrap(err, "importing existence columns") + } } // Import into fragment. - err = field.Import(req.RowIDs, req.ColumnIDs, timestamps) + err = field.Import(req.RowIDs, req.ColumnIDs, timestamps, opts...) if err != nil { api.server.logger.Printf("import error: index=%s, field=%s, shard=%d, columns=%d, err=%s", req.Index, req.Field, req.Shard, len(req.ColumnIDs), err) } diff --git a/client.go b/client.go index a2ca0bf96..889544557 100644 --- a/client.go +++ b/client.go @@ -37,8 +37,8 @@ type InternalClient interface { 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) error - ImportK(ctx context.Context, index, field string, bits []Bit) 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 EnsureField(ctx context.Context, indexName string, fieldName string) error EnsureFieldWithOptions(ctx context.Context, index, field string, opt FieldOptions) error @@ -103,10 +103,10 @@ 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) Import(ctx context.Context, index, field string, shard uint64, bits []Bit) error { +func (n nopInternalClient) Import(ctx context.Context, index, field string, shard uint64, bits []Bit, opts ...ImportOption) error { return nil } -func (n nopInternalClient) ImportK(ctx context.Context, index, field string, bits []Bit) error { +func (n nopInternalClient) ImportK(ctx context.Context, index, field string, bits []Bit, opts ...ImportOption) error { return nil } func (n nopInternalClient) ImportRoaring(ctx context.Context, uri *URI, index, field string, shard uint64, remote bool, data []byte) error { diff --git a/cmd/import.go b/cmd/import.go index 7bfca1a84..878aad286 100644 --- a/cmd/import.go +++ b/cmd/import.go @@ -63,6 +63,7 @@ omitted. If it is present then its format should be YYYY-MM-DDTHH:MM. flags.IntVarP(&Importer.BufferSize, "buffer-size", "s", 10000000, "Number of bits to buffer/sort before importing.") flags.BoolVarP(&Importer.Sort, "sort", "", false, "Enables sorting before import.") flags.BoolVarP(&Importer.CreateSchema, "create", "e", false, "Create the schema if it does not exist before import.") + flags.BoolVarP(&Importer.Clear, "clear", "", false, "Clear the data provided in the import.") ctl.SetTLSConfig(flags, &Importer.TLS.CertificatePath, &Importer.TLS.CertificateKeyPath, &Importer.TLS.SkipVerify) return importCmd diff --git a/cmd/import_test.go b/cmd/import_test.go index ba015dc04..7b566ce65 100644 --- a/cmd/import_test.go +++ b/cmd/import_test.go @@ -95,6 +95,17 @@ field = "f1" return v.Error() }, }, + { + args: []string{"import", "--index", "i1", "--field", "f1", "--clear", "true"}, + env: map[string]string{}, + validation: func() error { + v := validator{} + v.Check(cmd.Importer.Index, "i1") + v.Check(cmd.Importer.Field, "f1") + v.Check(cmd.Importer.Clear, true) + return v.Error() + }, + }, } executeDry(t, tests) } diff --git a/ctl/import.go b/ctl/import.go index d2694a64f..8f0f4a4e5 100644 --- a/ctl/import.go +++ b/ctl/import.go @@ -49,6 +49,9 @@ type ImportCommand struct { // nolint: maligned // CreateSchema ensures the schema exists before import CreateSchema bool + // Clear clears the import data as opposed to setting it. + Clear bool + // Filenames to import from. Paths []string `json:"paths"` @@ -255,7 +258,7 @@ func (cmd *ImportCommand) importBits(ctx context.Context, useColumnKeys, useRowK // If keys are used, all bits are sent to the primary translate store (i.e. coordinator). if useColumnKeys || useRowKeys { logger.Printf("importing keys: n=%d", len(bits)) - if err := cmd.client.ImportK(ctx, cmd.Index, cmd.Field, bits); err != nil { + if err := cmd.client.ImportK(ctx, cmd.Index, cmd.Field, bits, pilosa.OptImportOptionsClear(cmd.Clear)); err != nil { return errors.Wrap(err, "importing keys") } return nil @@ -272,7 +275,7 @@ func (cmd *ImportCommand) importBits(ctx context.Context, useColumnKeys, useRowK } logger.Printf("importing shard: %d, n=%d", shard, len(chunk)) - if err := cmd.client.Import(ctx, cmd.Index, cmd.Field, shard, chunk); err != nil { + if err := cmd.client.Import(ctx, cmd.Index, cmd.Field, shard, chunk, pilosa.OptImportOptionsClear(cmd.Clear)); err != nil { return errors.Wrap(err, "importing") } } diff --git a/ctl/import_test.go b/ctl/import_test.go index 2913005fb..fdf73720a 100644 --- a/ctl/import_test.go +++ b/ctl/import_test.go @@ -50,28 +50,55 @@ func TestImportCommand_Validation(t *testing.T) { } } -func TestImportCommand_Run(t *testing.T) { - buf := bytes.Buffer{} - stdin, stdout, stderr := GetIO(buf) - cm := NewImportCommand(stdin, stdout, stderr) - file, err := ioutil.TempFile("", "import.csv") - file.Write([]byte("1,2\n3,4\n5,6")) - ctx := context.Background() - if err != nil { - t.Fatal(err) - } +func TestImportCommand_Basic(t *testing.T) { + t.Run("set", func(t *testing.T) { + buf := bytes.Buffer{} + stdin, stdout, stderr := GetIO(buf) + cm := NewImportCommand(stdin, stdout, stderr) + file, err := ioutil.TempFile("", "import.csv") + file.Write([]byte("1,2\n3,4\n5,6")) + ctx := context.Background() + if err != nil { + t.Fatal(err) + } - cmd := test.MustRunCluster(t, 1)[0] - cm.Host = cmd.API.Node().URI.HostPort() + cmd := test.MustRunCluster(t, 1)[0] + cm.Host = cmd.API.Node().URI.HostPort() - cm.Index = "i" - cm.Field = "f" - cm.CreateSchema = true - cm.Paths = []string{file.Name()} - err = cm.Run(ctx) - if err != nil { - t.Fatalf("Import Run doesn't work: %s", err) - } + cm.Index = "i" + cm.Field = "f" + cm.CreateSchema = true + cm.Paths = []string{file.Name()} + err = cm.Run(ctx) + if err != nil { + t.Fatalf("Import Run doesn't work: %s", err) + } + }) + + t.Run("clear", func(t *testing.T) { + buf := bytes.Buffer{} + stdin, stdout, stderr := GetIO(buf) + cm := NewImportCommand(stdin, stdout, stderr) + file, err := ioutil.TempFile("", "import.csv") + file.Write([]byte("1,2\n3,4\n5,6")) + ctx := context.Background() + if err != nil { + t.Fatal(err) + } + + cmd := test.MustRunCluster(t, 1)[0] + cm.Host = cmd.API.Node().URI.HostPort() + + cm.Index = "i" + cm.Field = "f" + cm.CreateSchema = true + cm.Clear = true + cm.Paths = []string{file.Name()} + err = cm.Run(ctx) + if err != nil { + t.Fatalf("Import Run clear doesn't work: %s", err) + } + }) } // Ensure that the ImportValue path runs. diff --git a/field.go b/field.go index 63bb08b63..1a880a266 100644 --- a/field.go +++ b/field.go @@ -1047,11 +1047,25 @@ func (f *Field) Range(name string, op pql.Token, predicate int64) (*Row, error) } // Import bulk imports data. -func (f *Field) Import(rowIDs, columnIDs []uint64, timestamps []*time.Time) error { +func (f *Field) Import(rowIDs, columnIDs []uint64, timestamps []*time.Time, opts ...ImportOption) error { + + // Set up import options. + options := &ImportOptions{} + for _, opt := range opts { + err := opt(options) + if err != nil { + return errors.Wrap(err, "applying option") + } + } + // Determine quantum if timestamps are set. q := f.TimeQuantum() - if hasTime(timestamps) && q == "" { - return errors.New("time quantum not set in field") + if hasTime(timestamps) { + if q == "" { + return errors.New("time quantum not set in field") + } else if options.Clear { + return errors.New("import clear is not supported with timestamps") + } } fieldType := f.Type() @@ -1103,7 +1117,7 @@ func (f *Field) Import(rowIDs, columnIDs []uint64, timestamps []*time.Time) erro return errors.Wrap(err, "creating view") } - if err := frag.bulkImport(data.RowIDs, data.ColumnIDs); err != nil { + if err := frag.bulkImport(data.RowIDs, data.ColumnIDs, options); err != nil { return err } } diff --git a/fragment.go b/fragment.go index 63c9cce37..94e45ab1e 100644 --- a/fragment.go +++ b/fragment.go @@ -1418,20 +1418,20 @@ func (f *fragment) mergeBlock(id int, data []pairSet) (sets, clears []pairSet, e // bulkImport bulk imports a set of bits and then snapshots the storage. // The cache is updated to reflect the new data. -func (f *fragment) bulkImport(rowIDs, columnIDs []uint64) error { +func (f *fragment) bulkImport(rowIDs, columnIDs []uint64, options *ImportOptions) error { // Verify that there are an equal number of row ids and column ids. if len(rowIDs) != len(columnIDs) { return fmt.Errorf("mismatch of row/column len: %d != %d", len(rowIDs), len(columnIDs)) } - if f.mutexVector != nil { - return f.bulkImportMutex(rowIDs, columnIDs) + if f.mutexVector != nil && !options.Clear { + return f.bulkImportMutex(rowIDs, columnIDs, options) } - return f.bulkImportStandard(rowIDs, columnIDs) + return f.bulkImportStandard(rowIDs, columnIDs, options) } // bulkImportStandard performs a bulk import on a standard fragment. -func (f *fragment) bulkImportStandard(rowIDs, columnIDs []uint64) error { +func (f *fragment) bulkImportStandard(rowIDs, columnIDs []uint64, options *ImportOptions) error { // Create a temporary bitmap which will be populated by rowIDs and columnIDs // and then merged into the existing fragment's bitmap. localBitmap := roaring.NewBitmap() @@ -1480,10 +1480,18 @@ func (f *fragment) bulkImportStandard(rowIDs, columnIDs []uint64) error { // Merge localBitmap into fragment's existing data. var results *roaring.Bitmap - if f.storage.Count() > 0 { - results = f.storage.Union(localBitmap) + if options.Clear { + if f.storage.Count() > 0 { + results = f.storage.Difference(localBitmap) + } else { + results = roaring.NewBitmap() + } } else { - results = localBitmap + if f.storage.Count() > 0 { + results = f.storage.Union(localBitmap) + } else { + results = localBitmap + } } // Update cache counts for all affected rows. @@ -1500,7 +1508,7 @@ func (f *fragment) bulkImportStandard(rowIDs, columnIDs []uint64) error { // mutex restrictions. Because the mutex requirements must be checked // against storage, this method must acquire a write lock on the fragment // during the entire process, and it handles every bit independently. -func (f *fragment) bulkImportMutex(rowIDs, columnIDs []uint64) error { +func (f *fragment) bulkImportMutex(rowIDs, columnIDs []uint64, options *ImportOptions) error { f.mu.Lock() defer f.mu.Unlock() diff --git a/fragment_internal_test.go b/fragment_internal_test.go index be0db3b01..67ad70f45 100644 --- a/fragment_internal_test.go +++ b/fragment_internal_test.go @@ -1253,12 +1253,15 @@ func TestFragment_SetMutex(t *testing.T) { } } -// Ensure a fragment can import mutually exclusive values. -func TestFragment_ImportMutex(t *testing.T) { +// Ensure a fragment can import into set fields. +func TestFragment_ImportSet(t *testing.T) { tests := []struct { - rowIDs []uint64 - colIDs []uint64 - exp map[uint64][]uint64 + setRowIDs []uint64 + setColIDs []uint64 + setExp map[uint64][]uint64 + clearRowIDs []uint64 + clearColIDs []uint64 + clearExp map[uint64][]uint64 }{ { []uint64{1, 1, 1, 1}, @@ -1266,6 +1269,129 @@ func TestFragment_ImportMutex(t *testing.T) { map[uint64][]uint64{ 1: {0, 1, 2, 3}, }, + []uint64{}, + []uint64{}, + map[uint64][]uint64{ + 1: {0, 1, 2, 3}, + }, + }, + { + []uint64{1, 1, 1, 1, 2, 2, 2, 2}, + []uint64{0, 1, 2, 3, 0, 1, 2, 3}, + map[uint64][]uint64{ + 1: {0, 1, 2, 3}, + 2: {0, 1, 2, 3}, + }, + []uint64{1, 1, 2}, + []uint64{1, 2, 3}, + map[uint64][]uint64{ + 1: {0, 3}, + 2: {0, 1, 2}, + }, + }, + { + []uint64{1, 1, 1, 1, 2}, + []uint64{0, 1, 2, 3, 1}, + map[uint64][]uint64{ + 1: {0, 1, 2, 3}, + 2: {1}, + }, + []uint64{1, 1, 1, 1}, + []uint64{0, 1, 2, 3}, + map[uint64][]uint64{ + 1: {}, + 2: {1}, + }, + }, + { + []uint64{1, 1, 1, 1, 2, 2, 1}, + []uint64{0, 1, 2, 3, 1, 8, 1}, + map[uint64][]uint64{ + 1: {0, 1, 2, 3}, + 2: {1, 8}, + }, + []uint64{1, 1}, + []uint64{0, 0}, + map[uint64][]uint64{ + 1: {1, 2, 3}, + 2: {1, 8}, + }, + }, + { + []uint64{1, 2, 3}, + []uint64{8, 8, 8}, + map[uint64][]uint64{ + 1: {8}, + 2: {8}, + 3: {8}, + }, + []uint64{1, 2, 3}, + []uint64{9, 9, 9}, + map[uint64][]uint64{ + 1: {8}, + 2: {8}, + 3: {8}, + }, + }, + } + + for i, test := range tests { + t.Run(fmt.Sprintf("importset%d", i), func(t *testing.T) { + f := mustOpenFragment("i", "f", viewStandard, 0, "") + defer f.Close() + + // Set import. + err := f.bulkImport(test.setRowIDs, test.setColIDs, &ImportOptions{}) + if err != nil { + t.Fatalf("bulk importing ids: %v", err) + } + + // Check for expected results. + for k, v := range test.setExp { + cols := f.row(k).Columns() + if !reflect.DeepEqual(cols, v) { + t.Fatalf("expected: %v, but got: %v", v, cols) + } + } + + // Clear import. + err = f.bulkImport(test.clearRowIDs, test.clearColIDs, &ImportOptions{Clear: true}) + if err != nil { + t.Fatalf("bulk clearing ids: %v", err) + } + + // Check for expected results. + for k, v := range test.clearExp { + cols := f.row(k).Columns() + if !reflect.DeepEqual(cols, v) { + t.Fatalf("expected: %v, but got: %v", v, cols) + } + } + }) + } +} + +// Ensure a fragment can import mutually exclusive values. +func TestFragment_ImportMutex(t *testing.T) { + tests := []struct { + setRowIDs []uint64 + setColIDs []uint64 + setExp map[uint64][]uint64 + clearRowIDs []uint64 + clearColIDs []uint64 + clearExp map[uint64][]uint64 + }{ + { + []uint64{1, 1, 1, 1}, + []uint64{0, 1, 2, 3}, + map[uint64][]uint64{ + 1: {0, 1, 2, 3}, + }, + []uint64{}, + []uint64{}, + map[uint64][]uint64{ + 1: {0, 1, 2, 3}, + }, }, { []uint64{1, 1, 1, 1, 2, 2, 2, 2}, @@ -1274,6 +1400,12 @@ func TestFragment_ImportMutex(t *testing.T) { 1: {}, 2: {0, 1, 2, 3}, }, + []uint64{1, 1, 2}, + []uint64{1, 2, 3}, + map[uint64][]uint64{ + 1: {}, + 2: {0, 1, 2}, + }, }, { []uint64{1, 1, 1, 1, 2}, @@ -1282,6 +1414,12 @@ func TestFragment_ImportMutex(t *testing.T) { 1: {0, 2, 3}, 2: {1}, }, + []uint64{1, 1, 1, 1}, + []uint64{0, 1, 2, 3}, + map[uint64][]uint64{ + 1: {}, + 2: {1}, + }, }, { []uint64{1, 1, 1, 1, 2, 2, 1}, @@ -1290,6 +1428,12 @@ func TestFragment_ImportMutex(t *testing.T) { 1: {0, 1, 2, 3}, 2: {8}, }, + []uint64{1, 1}, + []uint64{0, 0}, + map[uint64][]uint64{ + 1: {1, 2, 3}, + 2: {8}, + }, }, { []uint64{1, 2, 3}, @@ -1299,6 +1443,13 @@ func TestFragment_ImportMutex(t *testing.T) { 2: {}, 3: {8}, }, + []uint64{1, 2, 3}, + []uint64{9, 9, 9}, + map[uint64][]uint64{ + 1: {}, + 2: {}, + 3: {8}, + }, }, } @@ -1307,13 +1458,28 @@ func TestFragment_ImportMutex(t *testing.T) { f := mustOpenMutexFragment("i", "f", viewStandard, 0, "") defer f.Close() - err := f.bulkImport(test.rowIDs, test.colIDs) + // Set import. + err := f.bulkImport(test.setRowIDs, test.setColIDs, &ImportOptions{}) if err != nil { t.Fatalf("bulk importing ids: %v", err) } // Check for expected results. - for k, v := range test.exp { + for k, v := range test.setExp { + cols := f.row(k).Columns() + if !reflect.DeepEqual(cols, v) { + t.Fatalf("expected: %v, but got: %v", v, cols) + } + } + + // Clear import. + err = f.bulkImport(test.clearRowIDs, test.clearColIDs, &ImportOptions{Clear: true}) + if err != nil { + t.Fatalf("bulk clearing ids: %v", err) + } + + // Check for expected results. + for k, v := range test.clearExp { cols := f.row(k).Columns() if !reflect.DeepEqual(cols, v) { t.Fatalf("expected: %v, but got: %v", v, cols) @@ -1326,9 +1492,12 @@ func TestFragment_ImportMutex(t *testing.T) { // Ensure a fragment can import bool values. func TestFragment_ImportBool(t *testing.T) { tests := []struct { - rowIDs []uint64 - colIDs []uint64 - exp map[uint64][]uint64 + setRowIDs []uint64 + setColIDs []uint64 + setExp map[uint64][]uint64 + clearRowIDs []uint64 + clearColIDs []uint64 + clearExp map[uint64][]uint64 }{ { []uint64{1, 1, 1, 1}, @@ -1336,6 +1505,11 @@ func TestFragment_ImportBool(t *testing.T) { map[uint64][]uint64{ 1: {0, 1, 2, 3}, }, + []uint64{}, + []uint64{}, + map[uint64][]uint64{ + 1: {0, 1, 2, 3}, + }, }, { []uint64{0, 0, 0, 0, 1, 1, 1, 1}, @@ -1344,6 +1518,13 @@ func TestFragment_ImportBool(t *testing.T) { 0: {}, 1: {0, 1, 2, 3}, }, + []uint64{1, 1, 2}, + []uint64{1, 2, 3}, + map[uint64][]uint64{ + 0: {}, + 1: {0, 3}, + 2: {}, + }, }, { []uint64{0, 0, 0, 0, 1}, @@ -1352,6 +1533,12 @@ func TestFragment_ImportBool(t *testing.T) { 0: {0, 2, 3}, 1: {1}, }, + []uint64{1, 1, 1, 1}, + []uint64{0, 1, 2, 3}, + map[uint64][]uint64{ + 0: {0, 2, 3}, + 1: {}, + }, }, { []uint64{1, 1, 1, 1, 0, 0, 1}, @@ -1360,6 +1547,12 @@ func TestFragment_ImportBool(t *testing.T) { 0: {8}, 1: {0, 1, 2, 3}, }, + []uint64{1, 1}, + []uint64{0, 0}, + map[uint64][]uint64{ + 0: {8}, + 1: {1, 2, 3}, + }, }, { []uint64{0, 1, 2}, @@ -1369,6 +1562,13 @@ func TestFragment_ImportBool(t *testing.T) { 1: {}, // This isn't {8} because fragment doesn't validate bool values. 2: {8}, }, + []uint64{1, 2, 3}, + []uint64{9, 9, 9}, + map[uint64][]uint64{ + 0: {}, + 1: {}, + 2: {8}, + }, }, } @@ -1377,13 +1577,28 @@ func TestFragment_ImportBool(t *testing.T) { f := mustOpenBoolFragment("i", "f", viewStandard, 0, "") defer f.Close() - err := f.bulkImport(test.rowIDs, test.colIDs) + // Set import. + err := f.bulkImport(test.setRowIDs, test.setColIDs, &ImportOptions{}) if err != nil { t.Fatalf("bulk importing ids: %v", err) } // Check for expected results. - for k, v := range test.exp { + for k, v := range test.setExp { + cols := f.row(k).Columns() + if !reflect.DeepEqual(cols, v) { + t.Fatalf("expected: %v, but got: %v", v, cols) + } + } + + // Clear import. + err = f.bulkImport(test.clearRowIDs, test.clearColIDs, &ImportOptions{Clear: true}) + if err != nil { + t.Fatalf("bulk importing ids: %v", err) + } + + // Check for expected results. + for k, v := range test.clearExp { cols := f.row(k).Columns() if !reflect.DeepEqual(cols, v) { t.Fatalf("expected: %v, but got: %v", v, cols) @@ -1427,6 +1642,7 @@ func BenchmarkFragment_FullSnapshot(b *testing.B) { rows := make([]uint64, sz) cols := make([]uint64, sz) + options := &ImportOptions{} max := 0 for row := 0; row < 100; row++ { val := 1 @@ -1437,7 +1653,7 @@ func BenchmarkFragment_FullSnapshot(b *testing.B) { val += 2 i++ } - if err := f.bulkImport(rows, cols); err != nil { + if err := f.bulkImport(rows, cols, options); err != nil { b.Fatalf("Error Building Sample: %s", err) } if row > max { @@ -1477,8 +1693,9 @@ func BenchmarkFragment_Import(b *testing.B) { } b.ResetTimer() b.ReportAllocs() + options := &ImportOptions{} for i := 0; i < b.N; i++ { - if err := f.bulkImport(rows, cols); err != nil { + if err := f.bulkImport(rows, cols, options); err != nil { b.Fatalf("Error Building Sample: %s", err) } } @@ -1692,7 +1909,8 @@ func TestFragment_RoaringImportTopN(t *testing.T) { f := mustOpenFragment("i", "f", viewStandard, 0, CacheTypeRanked) defer f.Close() - err := f.bulkImport(test.rowIDs, test.colIDs) + options := &ImportOptions{} + err := f.bulkImport(test.rowIDs, test.colIDs, options) if err != nil { t.Fatalf("bulk importing ids: %v", err) } @@ -1705,7 +1923,7 @@ func TestFragment_RoaringImportTopN(t *testing.T) { t.Fatalf("post bulk import:\n exp: %v\n got: %v\n", expPairs, pairs) } - err = f.bulkImport(test.rowIDs2, test.colIDs2) + err = f.bulkImport(test.rowIDs2, test.colIDs2, options) if err != nil { t.Fatalf("bulk importing ids: %v", err) } diff --git a/http/client.go b/http/client.go index 922ac6234..39f574bb9 100644 --- a/http/client.go +++ b/http/client.go @@ -294,13 +294,22 @@ func (c *InternalClient) QueryNode(ctx context.Context, uri *pilosa.URI, index s } // Import bulk imports bits for a single shard to a host. -func (c *InternalClient) Import(ctx context.Context, index, field string, shard uint64, bits []pilosa.Bit) error { +func (c *InternalClient) Import(ctx context.Context, index, field string, shard uint64, bits []pilosa.Bit, opts ...pilosa.ImportOption) error { if index == "" { return pilosa.ErrIndexRequired } else if field == "" { return pilosa.ErrFieldRequired } + // Set up import options. + options := &pilosa.ImportOptions{} + for _, opt := range opts { + err := opt(options) + if err != nil { + return errors.Wrap(err, "applying option") + } + } + buf, err := c.marshalImportPayload(index, field, shard, bits) if err != nil { return fmt.Errorf("Error Creating Payload: %s", err) @@ -314,7 +323,7 @@ func (c *InternalClient) Import(ctx context.Context, index, field string, shard // Import to each node. for _, node := range nodes { - if err := c.importNode(ctx, node, index, field, buf); err != nil { + if err := c.importNode(ctx, node, index, field, buf, options); err != nil { return fmt.Errorf("import node: host=%s, err=%s", node.URI, err) } } @@ -332,13 +341,22 @@ func getCoordinatorNode(nodes []*pilosa.Node) *pilosa.Node { } // ImportK bulk imports bits specified by string keys to a host. -func (c *InternalClient) ImportK(ctx context.Context, index, field string, bits []pilosa.Bit) error { +func (c *InternalClient) ImportK(ctx context.Context, index, field string, bits []pilosa.Bit, opts ...pilosa.ImportOption) error { if index == "" { return pilosa.ErrIndexRequired } else if field == "" { return pilosa.ErrFieldRequired } + // Set up import options. + options := &pilosa.ImportOptions{} + for _, opt := range opts { + err := opt(options) + if err != nil { + return errors.Wrap(err, "applying option") + } + } + buf, err := c.marshalImportPayload(index, field, 0, bits) if err != nil { return fmt.Errorf("Error Creating Payload: %s", err) @@ -356,7 +374,7 @@ func (c *InternalClient) ImportK(ctx context.Context, index, field string, bits } // Import to node. - if err := c.importNode(ctx, coord, index, field, buf); err != nil { + if err := c.importNode(ctx, coord, index, field, buf, options); err != nil { return fmt.Errorf("import node: host=%s, err=%s", coord.URI, err) } @@ -410,11 +428,17 @@ func (c *InternalClient) marshalImportPayload(index, field string, shard uint64, } // importNode sends a pre-marshaled import request to a node. -func (c *InternalClient) importNode(ctx context.Context, node *pilosa.Node, index, field string, buf []byte) error { +func (c *InternalClient) importNode(ctx context.Context, node *pilosa.Node, index, field string, buf []byte, opts *pilosa.ImportOptions) error { // Create URL & HTTP request. path := fmt.Sprintf("/index/%s/field/%s/import", index, field) u := nodePathToURL(node, path) - req, err := http.NewRequest("POST", u.String(), bytes.NewReader(buf)) + + url := u.String() + if opts.Clear { + url += "?clear=true" + } + + req, err := http.NewRequest("POST", url, bytes.NewReader(buf)) if err != nil { return errors.Wrap(err, "creating request") } @@ -467,9 +491,12 @@ func (c *InternalClient) ImportValue(ctx context.Context, index, field string, s return fmt.Errorf("shard nodes: %s", err) } + // Set up import options. + options := &pilosa.ImportOptions{} + // Import to each node. for _, node := range nodes { - if err := c.importNode(ctx, node, index, field, buf); err != nil { + if err := c.importNode(ctx, node, index, field, buf, options); err != nil { return fmt.Errorf("import node: host=%s, err=%s", node.URI, err) } } @@ -495,8 +522,11 @@ func (c *InternalClient) ImportValueK(ctx context.Context, index, field string, return fmt.Errorf("could not find the coordinator node") } + // Set up import options. + options := &pilosa.ImportOptions{} + // Import to node. - if err := c.importNode(ctx, coord, index, field, buf); err != nil { + if err := c.importNode(ctx, coord, index, field, buf, options); err != nil { return fmt.Errorf("import node: host=%s, err=%s", coord.URI, err) } diff --git a/http/handler.go b/http/handler.go index 8b239ec2d..bfa3d5e4a 100644 --- a/http/handler.go +++ b/http/handler.go @@ -181,8 +181,8 @@ func (h *Handler) populateValidators() { h.validators["DeleteIndex"] = queryValidationSpecRequired() h.validators["PostField"] = queryValidationSpecRequired() h.validators["DeleteField"] = queryValidationSpecRequired() - h.validators["PostImport"] = queryValidationSpecRequired() - h.validators["PostImportRoaring"] = queryValidationSpecRequired().Optional("remote") + h.validators["PostImport"] = queryValidationSpecRequired().Optional("clear") + h.validators["PostImportRoaring"] = queryValidationSpecRequired().Optional("remote", "clear") h.validators["PostQuery"] = queryValidationSpecRequired().Optional("shards", "columnAttrs", "excludeRowAttrs", "excludeColumns") h.validators["GetInfo"] = queryValidationSpecRequired() h.validators["RecalculateCaches"] = queryValidationSpecRequired() @@ -994,6 +994,10 @@ func (h *Handler) handlePostImport(w http.ResponseWriter, r *http.Request) { indexName := mux.Vars(r)["index"] fieldName := mux.Vars(r)["field"] + // If the clear flag is true, treat the import as clear bits. + q := r.URL.Query() + doClear := q.Get("clear") == "true" + // Get index and field type to determine how to handle the // import data. field, err := h.api.Field(r.Context(), indexName, fieldName) @@ -1044,7 +1048,7 @@ func (h *Handler) handlePostImport(w http.ResponseWriter, r *http.Request) { return } - if err := h.api.Import(r.Context(), req); err != nil { + if err := h.api.Import(r.Context(), req, pilosa.OptImportOptionsClear(doClear)); err != nil { switch errors.Cause(err) { case pilosa.ErrClusterDoesNotOwnShard: http.Error(w, err.Error(), http.StatusPreconditionFailed) From a7a15c64a21ff2d54c6d4bd7b8e6feb4ddaeff23 Mon Sep 17 00:00:00 2001 From: Travis Turner Date: Mon, 22 Oct 2018 15:59:24 -0500 Subject: [PATCH 06/11] support clear imports to int fields. fix bug in fragment.sum --- api.go | 21 +++++++++--- client.go | 8 ++--- ctl/import.go | 2 +- ctl/import_test.go | 69 +++++++++++++++++++++++++++------------ field.go | 4 +-- fragment.go | 69 +++++++++++++++++++++++++-------------- fragment_internal_test.go | 28 ++++++++++++++++ http/client.go | 28 +++++++++++----- http/handler.go | 2 +- 9 files changed, 165 insertions(+), 66 deletions(-) diff --git a/api.go b/api.go index bb900db4e..a81f1db12 100644 --- a/api.go +++ b/api.go @@ -757,11 +757,20 @@ func (api *API) Import(_ context.Context, req *ImportRequest, opts ...ImportOpti } // ImportValue bulk imports values into a particular field. -func (api *API) ImportValue(_ context.Context, req *ImportValueRequest) error { +func (api *API) ImportValue(_ context.Context, req *ImportValueRequest, opts ...ImportOption) error { if err := api.validate(apiImportValue); err != nil { return errors.Wrap(err, "validating api method") } + // Set up import options. + options := &ImportOptions{} + for _, opt := range opts { + err := opt(options) + if err != nil { + return errors.Wrap(err, "applying option") + } + } + index := api.holder.Index(req.Index) if index == nil { return newNotFoundError(ErrIndexNotFound) @@ -783,13 +792,15 @@ func (api *API) ImportValue(_ context.Context, req *ImportValueRequest) error { } // Import columnIDs into existence field. - if err := importExistenceColumns(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 errors.Wrap(err, "importing existence columns") + if !options.Clear { + if err := importExistenceColumns(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 errors.Wrap(err, "importing existence columns") + } } // Import into fragment. - err = field.importValue(req.ColumnIDs, req.Values) + err = field.importValue(req.ColumnIDs, req.Values, options) if err != nil { api.server.logger.Printf("import error: index=%s, field=%s, shard=%d, columns=%d, err=%s", req.Index, req.Field, req.Shard, len(req.ColumnIDs), err) } diff --git a/client.go b/client.go index 889544557..e6a02aa96 100644 --- a/client.go +++ b/client.go @@ -42,8 +42,8 @@ type InternalClient interface { EnsureIndex(ctx context.Context, name string, options IndexOptions) error EnsureField(ctx context.Context, indexName string, fieldName string) error EnsureFieldWithOptions(ctx context.Context, index, field string, opt FieldOptions) error - ImportValue(ctx context.Context, index, field string, shard uint64, vals []FieldValue) error - ImportValueK(ctx context.Context, index, field string, vals []FieldValue) error + ImportValue(ctx context.Context, index, field string, shard uint64, vals []FieldValue, opts ...ImportOption) error + ImportValueK(ctx context.Context, index, field string, vals []FieldValue, opts ...ImportOption) error ExportCSV(ctx context.Context, index, field string, shard uint64, w io.Writer) error CreateField(ctx context.Context, index, field string) error CreateFieldWithOptions(ctx context.Context, index, field string, opt FieldOptions) error @@ -121,10 +121,10 @@ func (n nopInternalClient) EnsureField(ctx context.Context, indexName string, fi func (n nopInternalClient) EnsureFieldWithOptions(ctx context.Context, index, field string, opt FieldOptions) error { return nil } -func (n nopInternalClient) ImportValue(ctx context.Context, index, field string, shard uint64, vals []FieldValue) error { +func (n nopInternalClient) ImportValue(ctx context.Context, index, field string, shard uint64, vals []FieldValue, opts ...ImportOption) error { return nil } -func (n nopInternalClient) ImportValueK(ctx context.Context, index, field string, vals []FieldValue) error { +func (n nopInternalClient) ImportValueK(ctx context.Context, index, field string, vals []FieldValue, opts ...ImportOption) error { return nil } func (n nopInternalClient) ExportCSV(ctx context.Context, index, field string, shard uint64, w io.Writer) error { diff --git a/ctl/import.go b/ctl/import.go index 8f0f4a4e5..d49b50eb3 100644 --- a/ctl/import.go +++ b/ctl/import.go @@ -380,7 +380,7 @@ func (cmd *ImportCommand) importValues(ctx context.Context, useColumnKeys bool, } logger.Printf("importing shard: %d, n=%d", shard, len(vals)) - if err := cmd.client.ImportValue(ctx, cmd.Index, cmd.Field, shard, vals); err != nil { + if err := cmd.client.ImportValue(ctx, cmd.Index, cmd.Field, shard, vals, pilosa.OptImportOptionsClear(cmd.Clear)); err != nil { return errors.Wrap(err, "importing values") } } diff --git a/ctl/import_test.go b/ctl/import_test.go index fdf73720a..ad4fd6cf0 100644 --- a/ctl/import_test.go +++ b/ctl/import_test.go @@ -103,29 +103,58 @@ func TestImportCommand_Basic(t *testing.T) { // Ensure that the ImportValue path runs. func TestImportCommand_RunValue(t *testing.T) { - buf := bytes.Buffer{} - stdin, stdout, stderr := GetIO(buf) - cm := NewImportCommand(stdin, stdout, stderr) - file, err := ioutil.TempFile("", "import-value.csv") - file.Write([]byte("1,2\n3,4\n5,6")) - ctx := context.Background() - if err != nil { - t.Fatal(err) - } + t.Run("set", func(t *testing.T) { + buf := bytes.Buffer{} + stdin, stdout, stderr := GetIO(buf) + cm := NewImportCommand(stdin, stdout, stderr) + file, err := ioutil.TempFile("", "import-value.csv") + file.Write([]byte("1,2\n3,4\n5,6")) + ctx := context.Background() + if err != nil { + t.Fatal(err) + } - cmd := test.MustRunCluster(t, 1)[0] - cm.Host = cmd.API.Node().URI.HostPort() + cmd := test.MustRunCluster(t, 1)[0] + cm.Host = cmd.API.Node().URI.HostPort() - http.DefaultClient.Do(MustNewHTTPRequest("POST", "http://"+cm.Host+"/index/i", strings.NewReader(""))) - http.DefaultClient.Do(MustNewHTTPRequest("POST", "http://"+cm.Host+"/index/i/field/f", strings.NewReader(`{"options":{"type": "int", "min": 0, "max": 100}}`))) + http.DefaultClient.Do(MustNewHTTPRequest("POST", "http://"+cm.Host+"/index/i", strings.NewReader(""))) + http.DefaultClient.Do(MustNewHTTPRequest("POST", "http://"+cm.Host+"/index/i/field/f", strings.NewReader(`{"options":{"type": "int", "min": 0, "max": 100}}`))) - cm.Index = "i" - cm.Field = "f" - cm.Paths = []string{file.Name()} - err = cm.Run(ctx) - if err != nil { - t.Fatalf("Import Run with values doesn't work: %s", err) - } + cm.Index = "i" + cm.Field = "f" + cm.Paths = []string{file.Name()} + err = cm.Run(ctx) + if err != nil { + t.Fatalf("Import Run with values doesn't work: %s", err) + } + }) + + t.Run("clear", func(t *testing.T) { + buf := bytes.Buffer{} + stdin, stdout, stderr := GetIO(buf) + cm := NewImportCommand(stdin, stdout, stderr) + file, err := ioutil.TempFile("", "import-value.csv") + file.Write([]byte("1,2\n3,4\n5,6")) + ctx := context.Background() + if err != nil { + t.Fatal(err) + } + + cmd := test.MustRunCluster(t, 1)[0] + cm.Host = cmd.API.Node().URI.HostPort() + + http.DefaultClient.Do(MustNewHTTPRequest("POST", "http://"+cm.Host+"/index/i", strings.NewReader(""))) + http.DefaultClient.Do(MustNewHTTPRequest("POST", "http://"+cm.Host+"/index/i/field/f", strings.NewReader(`{"options":{"type": "int", "min": 0, "max": 100}}`))) + + cm.Index = "i" + cm.Field = "f" + cm.Paths = []string{file.Name()} + cm.Clear = true + err = cm.Run(ctx) + if err != nil { + t.Fatalf("Import Run with values doesn't work: %s", err) + } + }) } // Ensure that import with keys runs. diff --git a/field.go b/field.go index 1a880a266..3422e0e19 100644 --- a/field.go +++ b/field.go @@ -1126,7 +1126,7 @@ func (f *Field) Import(rowIDs, columnIDs []uint64, timestamps []*time.Time, opts } // importValue bulk imports range-encoded value data. -func (f *Field) importValue(columnIDs []uint64, values []int64) error { +func (f *Field) importValue(columnIDs []uint64, values []int64, options *ImportOptions) error { viewName := viewBSIGroupPrefix + f.name // Get the bsiGroup so we know bitDepth. bsig := f.bsiGroup(f.name) @@ -1174,7 +1174,7 @@ func (f *Field) importValue(columnIDs []uint64, values []int64) error { baseValues[i] = uint64(value - bsig.Min) } - if err := frag.importValue(data.ColumnIDs, baseValues, bsig.BitDepth()); err != nil { + if err := frag.importValue(data.ColumnIDs, baseValues, bsig.BitDepth(), options.Clear); err != nil { return err } } diff --git a/fragment.go b/fragment.go index 94e45ab1e..d45757762 100644 --- a/fragment.go +++ b/fragment.go @@ -612,8 +612,17 @@ func (f *fragment) value(columnID uint64, bitDepth uint) (value uint64, exists b return value, true, nil } +// clearValue uses a column of bits to clear a multi-bit value. +func (f *fragment) clearValue(columnID uint64, bitDepth uint, value uint64) (changed bool, err error) { + return f.setValueBase(columnID, bitDepth, value, true) +} + // setValue uses a column of bits to set a multi-bit value. func (f *fragment) setValue(columnID uint64, bitDepth uint, value uint64) (changed bool, err error) { + return f.setValueBase(columnID, bitDepth, value, false) +} + +func (f *fragment) setValueBase(columnID uint64, bitDepth uint, value uint64, clear bool) (changed bool, err error) { f.mu.Lock() defer f.mu.Unlock() @@ -633,19 +642,26 @@ func (f *fragment) setValue(columnID uint64, bitDepth uint, value uint64) (chang } } - // Mark value as set. - if c, err := f.unprotectedSetBit(uint64(bitDepth), columnID); err != nil { - return changed, errors.Wrap(err, "marking not-null") - } else if c { - changed = true + // Mark value as set (or cleared). + if clear { + if c, err := f.unprotectedClearBit(uint64(bitDepth), columnID); err != nil { + return changed, errors.Wrap(err, "clearing not-null") + } else if c { + changed = true + } + } else { + if c, err := f.unprotectedSetBit(uint64(bitDepth), columnID); err != nil { + return changed, errors.Wrap(err, "marking not-null") + } else if c { + changed = true + } } return changed, nil } // importSetValue is a more efficient SetValue just for imports. -func (f *fragment) importSetValue(columnID uint64, bitDepth uint, value uint64) (changed bool, err error) { // nolint: unparam - +func (f *fragment) importSetValue(columnID uint64, bitDepth uint, value uint64, clear bool) (changed bool, err error) { // nolint: unparam for i := uint(0); i < bitDepth; i++ { if value&(1< Date: Mon, 22 Oct 2018 16:25:32 -0500 Subject: [PATCH 07/11] add clear support for ImportRoaring --- api.go | 16 +++++++++++++--- client.go | 4 ++-- field.go | 4 ++-- fragment.go | 10 +++++++--- fragment_internal_test.go | 4 ++-- http/client.go | 14 +++++++++++++- 6 files changed, 39 insertions(+), 13 deletions(-) diff --git a/api.go b/api.go index a81f1db12..df1265433 100644 --- a/api.go +++ b/api.go @@ -256,10 +256,20 @@ func (api *API) Field(_ context.Context, indexName, fieldName string) (*Field, e // (shard*ShardWidth)+(i%ShardWidth). That is to say that "data" represents all // of the rows in this shard of this field concatenated together in one long // bitmap. -func (api *API) ImportRoaring(ctx context.Context, indexName, fieldName string, shard uint64, remote bool, data []byte) (err error) { +func (api *API) ImportRoaring(ctx context.Context, indexName, fieldName string, shard uint64, remote bool, data []byte, opts ...ImportOption) (err error) { if err = api.validate(apiField); err != nil { return errors.Wrap(err, "validating api method") } + + // Set up import options. + options := &ImportOptions{} + for _, opt := range opts { + err := opt(options) + if err != nil { + return errors.Wrap(err, "applying option") + } + } + nodes := api.cluster.shardNodes(indexName, shard) var eg errgroup.Group @@ -280,14 +290,14 @@ func (api *API) ImportRoaring(ctx context.Context, indexName, fieldName string, d2 := make([]byte, len(data)) copy(d2, data) eg.Go(func() error { - return field.importRoaring(d2, shard) + return field.importRoaring(d2, shard, options.Clear) }) go func(node *Node) { }(node) } else if !remote { // if remote == true we don't forward to other nodes // forward it on eg.Go(func() error { - return api.server.defaultClient.ImportRoaring(ctx, &node.URI, indexName, fieldName, shard, true, data) + return api.server.defaultClient.ImportRoaring(ctx, &node.URI, indexName, fieldName, shard, true, data, opts...) }) } } diff --git a/client.go b/client.go index e6a02aa96..f6c49d14b 100644 --- a/client.go +++ b/client.go @@ -53,7 +53,7 @@ type InternalClient interface { RowAttrDiff(ctx context.Context, uri *URI, index, field string, blks []AttrBlock) (map[uint64]map[string]interface{}, error) SendMessage(ctx context.Context, uri *URI, msg []byte) error RetrieveShardFromURI(ctx context.Context, index, field string, shard uint64, uri URI) (io.ReadCloser, error) - ImportRoaring(ctx context.Context, uri *URI, index, field string, shard uint64, remote bool, data []byte) error + ImportRoaring(ctx context.Context, uri *URI, index, field string, shard uint64, remote bool, data []byte, opts ...ImportOption) error } //=============== @@ -109,7 +109,7 @@ func (n nopInternalClient) Import(ctx context.Context, index, field string, shar func (n nopInternalClient) ImportK(ctx context.Context, index, field string, bits []Bit, opts ...ImportOption) error { return nil } -func (n nopInternalClient) ImportRoaring(ctx context.Context, uri *URI, index, field string, shard uint64, remote bool, data []byte) error { +func (n nopInternalClient) ImportRoaring(ctx context.Context, uri *URI, index, field string, shard uint64, remote bool, data []byte, opts ...ImportOption) error { return nil } func (n nopInternalClient) EnsureIndex(ctx context.Context, name string, options IndexOptions) error { diff --git a/field.go b/field.go index 3422e0e19..bc656fa54 100644 --- a/field.go +++ b/field.go @@ -1182,7 +1182,7 @@ func (f *Field) importValue(columnIDs []uint64, values []int64, options *ImportO return nil } -func (f *Field) importRoaring(data []byte, shard uint64) error { +func (f *Field) importRoaring(data []byte, shard uint64, clear bool) error { viewName := viewStandard view, err := f.createViewIfNotExists(viewName) @@ -1195,7 +1195,7 @@ func (f *Field) importRoaring(data []byte, shard uint64) error { return errors.Wrap(err, "creating fragment") } - if err := frag.importRoaring(data); err != nil { + if err := frag.importRoaring(data, clear); err != nil { return err } diff --git a/fragment.go b/fragment.go index d45757762..10e5e3920 100644 --- a/fragment.go +++ b/fragment.go @@ -1651,7 +1651,7 @@ func (f *fragment) importValue(columnIDs, values []uint64, bitDepth uint, clear // importRoaring imports from the official roaring data format defined at // https://github.com/RoaringBitmap/RoaringFormatSpec or from pilosa's version // of the roaring format. The cache is updated to reflect the new data. -func (f *fragment) importRoaring(data []byte) error { +func (f *fragment) importRoaring(data []byte, clear bool) error { f.mu.Lock() defer f.mu.Unlock() bm := roaring.NewBitmap() @@ -1679,8 +1679,12 @@ func (f *fragment) importRoaring(data []byte) error { lastRow = vRow } - if f.storage.Count() > 0 { - bm = f.storage.Union(bm) + if clear { + bm = f.storage.Difference(bm) + } else { + if f.storage.Count() > 0 { + bm = f.storage.Union(bm) + } } for _, rowID := range rowSet { diff --git a/fragment_internal_test.go b/fragment_internal_test.go index 5097fb680..ac080b5b7 100644 --- a/fragment_internal_test.go +++ b/fragment_internal_test.go @@ -1900,7 +1900,7 @@ func TestFragment_RoaringImport(t *testing.T) { if err != nil { t.Fatalf("writing to buffer: %v", err) } - f.importRoaring(buf.Bytes()) + f.importRoaring(buf.Bytes(), false) exp := calcExpected(test[:num+1]...) for row, expCols := range exp { cols := f.row(uint64(row)).Columns() @@ -1972,7 +1972,7 @@ func TestFragment_RoaringImportTopN(t *testing.T) { if err != nil { t.Fatalf("writing to buffer: %v", err) } - f.importRoaring(buf.Bytes()) + f.importRoaring(buf.Bytes(), false) rows, cols := toRowsCols(test.roaring) expPairs = calcTop(append(test.rowIDs, rows...), append(test.colIDs, cols...)) pairs, err = f.top(topOptions{}) diff --git a/http/client.go b/http/client.go index d546c5b9c..17c0b2d63 100644 --- a/http/client.go +++ b/http/client.go @@ -569,7 +569,7 @@ func (c *InternalClient) marshalImportValuePayload(index, field string, shard ui // ImportRoaring does fast import of raw bits in roaring format (pilosa or // official format, see API.ImportRoaring). -func (c *InternalClient) ImportRoaring(ctx context.Context, uri *pilosa.URI, index, field string, shard uint64, remote bool, data []byte) error { +func (c *InternalClient) ImportRoaring(ctx context.Context, uri *pilosa.URI, index, field string, shard uint64, remote bool, data []byte, opts ...pilosa.ImportOption) error { if index == "" { return pilosa.ErrIndexRequired } else if field == "" { @@ -579,7 +579,19 @@ func (c *InternalClient) ImportRoaring(ctx context.Context, uri *pilosa.URI, ind uri = c.defaultURI } + // Set up import options. + options := &pilosa.ImportOptions{} + for _, opt := range opts { + err := opt(options) + if err != nil { + return errors.Wrap(err, "applying option") + } + } + url := fmt.Sprintf("%s/index/%s/field/%s/import-roaring/%d?remote=%v", uri, index, field, shard, remote) + if options.Clear { + url += "&clear=true" + } // Generate HTTP request. req, err := http.NewRequest("POST", url, bytes.NewBuffer(data)) From f29b6e79ff1856dd3e7661302b9119559262014e Mon Sep 17 00:00:00 2001 From: Travis Turner Date: Mon, 22 Oct 2018 16:44:56 -0500 Subject: [PATCH 08/11] add import clear test coverage to http package --- http/client_test.go | 86 +++++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 86 insertions(+) diff --git a/http/client_test.go b/http/client_test.go index accbf86e4..be6bdb996 100644 --- a/http/client_test.go +++ b/http/client_test.go @@ -362,6 +362,22 @@ func TestClient_Import(t *testing.T) { if a := hldr.Row("i", "f", 200).Columns(); !reflect.DeepEqual(a, []uint64{6}) { t.Fatalf("unexpected columns: %+v", a) } + + // Clear some data. + if err := c.Import(context.Background(), "i", "f", 0, []pilosa.Bit{ + {RowID: 0, ColumnID: 5}, + {RowID: 200, ColumnID: 6}, + }, pilosa.OptImportOptionsClear(true)); err != nil { + t.Fatal(err) + } + + // Verify data. + if a := hldr.Row("i", "f", 0).Columns(); !reflect.DeepEqual(a, []uint64{1}) { + t.Fatalf("unexpected columns: %+v", a) + } + if a := hldr.Row("i", "f", 200).Columns(); !reflect.DeepEqual(a, []uint64{}) { + t.Fatalf("unexpected columns: %+v", a) + } } // Ensure client can bulk import data. @@ -638,6 +654,37 @@ func TestClient_ImportKeys(t *testing.T) { if !reflect.DeepEqual(result.Results[0].(*pilosa.Row).Keys, []string{"col2", "col3"}) { t.Fatalf("unexpected column keys: %s", spew.Sdump(result)) } + + // Clear data. + if err := c.ImportValue(context.Background(), "i", "f", 0, []pilosa.FieldValue{ + {ColumnKey: "col2", Value: 20}, + }, pilosa.OptImportOptionsClear(true)); err != nil { + t.Fatal(err) + } + + // Verify Sum. + sum, cnt, err = field.Sum(nil, fldName) + if err != nil { + t.Fatal(err) + } + if sum != 30 || cnt != 2 { + t.Fatalf("unexpected values: got sum=%v, count=%v; expected sum=30, cnt=2", sum, cnt) + } + + // Verify Range + queryRequest = &pilosa.QueryRequest{ + Query: fmt.Sprintf(`Range(%s>10)`, fldName), + Remote: false, + } + + result, err = c.Query(context.Background(), "i", queryRequest) + if err != nil { + t.Fatal(err) + } + + if !reflect.DeepEqual(result.Results[0].(*pilosa.Row).Keys, []string{"col3"}) { + t.Fatalf("unexpected column keys: %s", spew.Sdump(result)) + } }) } @@ -706,6 +753,45 @@ func TestClient_ImportValue(t *testing.T) { if max != 40 || cnt != 1 { t.Fatalf("unexpected values: got max=%v, count=%v; expected max=40, cnt=1", max, cnt) } + + // Send import request. + if err := c.ImportValue(context.Background(), "i", "f", 0, []pilosa.FieldValue{ + {ColumnID: 1, Value: -10}, + {ColumnID: 3, Value: 40}, + }, pilosa.OptImportOptionsClear(true)); err != nil { + t.Fatal(err) + } + + // Verify Sum. + sum, cnt, err = field.Sum(nil, fldName) + if err != nil { + t.Fatal(err) + } + if sum != 20 || cnt != 1 { + t.Fatalf("unexpected values: got sum=%v, count=%v; expected sum=20, cnt=1", sum, cnt) + } + + // Verify Min with Filter. + filter, err = field.Range(fldName, pql.GT, 40) + if err != nil { + t.Fatal(err) + } + min, cnt, err = field.Min(filter, fldName) + if err != nil { + t.Fatal(err) + } + if min != -100 || cnt != 0 { + t.Fatalf("unexpected values: got min=%v, count=%v; expected min=-100, cnt=0", min, cnt) + } + + // Verify Max. + max, cnt, err = field.Max(nil, fldName) + if err != nil { + t.Fatal(err) + } + if max != 20 || cnt != 1 { + t.Fatalf("unexpected values: got max=%v, count=%v; expected max=20, cnt=1", max, cnt) + } } // Ensure client can bulk import data while tracking existence. From 318e588b4dce4477b27379d8e379a401a06ec90e Mon Sep 17 00:00:00 2001 From: Travis Turner Date: Mon, 22 Oct 2018 17:05:18 -0500 Subject: [PATCH 09/11] document the --clear flag for imports --- docs/administration.md | 12 ++++++++++++ 1 file changed, 12 insertions(+) diff --git a/docs/administration.md b/docs/administration.md index e626cc0e5..3f07a5b68 100644 --- a/docs/administration.md +++ b/docs/administration.md @@ -82,6 +82,18 @@ For example, importing a file with the following contents will result in columns

Note that you must first create a field. View Create Field for more details. The `-e` flag can create the necessary schema when using a field of type "set".

+#### Clearing Data via Import + +By using the `--clear` flag with the import command, Pilosa will clear the values provided in the import payload. + +For example, importing a file with the following contents along with the `--clear` flag will result in data being cleared from row 0, column 9; row 1, columns 2 and 8; and row 3, column 12. Clearing a value that doesn't exists is allowed. +``` +0,9 +1,2 +1,8 +3,12 +``` + #### Exporting Exporting data to csv can be performed on a live instance of Pilosa. You need to specify the index and the field. The API also expects the shard number, but the `pilosa export` sub command will export all shards within a field. The data will be in csv format `Row,Column` and sorted by column. From f02605a52821dabbe8d246d2529f4b567579ab90 Mon Sep 17 00:00:00 2001 From: Travis Turner Date: Tue, 23 Oct 2018 17:35:17 -0500 Subject: [PATCH 10/11] consolidate ImportOptions setup. use url.Values{}. fix comments. --- api.go | 40 +++++++++++++++++++++------------------- http/client.go | 11 +++++++---- http/client_test.go | 4 ++-- 3 files changed, 30 insertions(+), 25 deletions(-) diff --git a/api.go b/api.go index df1265433..af206496b 100644 --- a/api.go +++ b/api.go @@ -239,6 +239,17 @@ func (api *API) Field(_ context.Context, indexName, fieldName string) (*Field, e return field, nil } +func setUpImportOptions(opts ...ImportOption) (*ImportOptions, error) { + options := &ImportOptions{} + for _, opt := range opts { + err := opt(options) + if err != nil { + return nil, errors.Wrap(err, "applying option") + } + } + return options, nil +} + // ImportRoaring is a low level interface for importing data to Pilosa when // extremely high throughput is desired. The data must be encoded in a // particular way which may be unintuitive (discussed below). The data is merged @@ -262,12 +273,9 @@ func (api *API) ImportRoaring(ctx context.Context, indexName, fieldName string, } // Set up import options. - options := &ImportOptions{} - for _, opt := range opts { - err := opt(options) - if err != nil { - return errors.Wrap(err, "applying option") - } + options, err := setUpImportOptions(opts...) + if err != nil { + return errors.Wrap(err, "setting up import options") } nodes := api.cluster.shardNodes(indexName, shard) @@ -685,7 +693,7 @@ type ImportOptions struct { Clear bool } -// ImportOption is a functional option type for API.Import +// ImportOption is a functional option type for API.Import. type ImportOption func(*ImportOptions) error func OptImportOptionsClear(c bool) ImportOption { @@ -702,12 +710,9 @@ func (api *API) Import(_ context.Context, req *ImportRequest, opts ...ImportOpti } // Set up import options. - options := &ImportOptions{} - for _, opt := range opts { - err := opt(options) - if err != nil { - return errors.Wrap(err, "applying option") - } + options, err := setUpImportOptions(opts...) + if err != nil { + return errors.Wrap(err, "setting up import options") } index := api.holder.Index(req.Index) @@ -773,12 +778,9 @@ func (api *API) ImportValue(_ context.Context, req *ImportValueRequest, opts ... } // Set up import options. - options := &ImportOptions{} - for _, opt := range opts { - err := opt(options) - if err != nil { - return errors.Wrap(err, "applying option") - } + options, err := setUpImportOptions(opts...) + if err != nil { + return errors.Wrap(err, "setting up import options") } index := api.holder.Index(req.Index) diff --git a/http/client.go b/http/client.go index 17c0b2d63..06bb31479 100644 --- a/http/client.go +++ b/http/client.go @@ -433,10 +433,11 @@ func (c *InternalClient) importNode(ctx context.Context, node *pilosa.Node, inde path := fmt.Sprintf("/index/%s/field/%s/import", index, field) u := nodePathToURL(node, path) - url := u.String() + vals := url.Values{} if opts.Clear { - url += "?clear=true" + vals.Set("clear", "true") } + url := fmt.Sprintf("%s?%s", u.String(), vals.Encode()) req, err := http.NewRequest("POST", url, bytes.NewReader(buf)) if err != nil { @@ -588,10 +589,12 @@ func (c *InternalClient) ImportRoaring(ctx context.Context, uri *pilosa.URI, ind } } - url := fmt.Sprintf("%s/index/%s/field/%s/import-roaring/%d?remote=%v", uri, index, field, shard, remote) + vals := url.Values{} + vals.Set("remote", strconv.FormatBool(remote)) if options.Clear { - url += "&clear=true" + vals.Set("clear", "true") } + url := fmt.Sprintf("%s/index/%s/field/%s/import-roaring/%d?%s", uri, index, field, shard, vals.Encode()) // Generate HTTP request. req, err := http.NewRequest("POST", url, bytes.NewBuffer(data)) diff --git a/http/client_test.go b/http/client_test.go index be6bdb996..4ba076408 100644 --- a/http/client_test.go +++ b/http/client_test.go @@ -640,7 +640,7 @@ func TestClient_ImportKeys(t *testing.T) { t.Fatalf("unexpected values: got sum=%v, count=%v; expected sum=50, cnt=3", sum, cnt) } - // Verify Range + // Verify Range. queryRequest := &pilosa.QueryRequest{ Query: fmt.Sprintf(`Range(%s>10)`, fldName), Remote: false, @@ -671,7 +671,7 @@ func TestClient_ImportKeys(t *testing.T) { t.Fatalf("unexpected values: got sum=%v, count=%v; expected sum=30, cnt=2", sum, cnt) } - // Verify Range + // Verify Range. queryRequest = &pilosa.QueryRequest{ Query: fmt.Sprintf(`Range(%s>10)`, fldName), Remote: false, From e69a8f3db4e9971090fdac3454beeccbfbff725d Mon Sep 17 00:00:00 2001 From: Travis Turner Date: Wed, 24 Oct 2018 15:16:08 -0500 Subject: [PATCH 11/11] handle clear flag on import roaring (http package) --- http/client_test.go | 72 +++++++++++++++++++++++++++++++++++++++++++-- http/handler.go | 5 +++- 2 files changed, 73 insertions(+), 4 deletions(-) diff --git a/http/client_test.go b/http/client_test.go index 4ba076408..23d8b8573 100644 --- a/http/client_test.go +++ b/http/client_test.go @@ -408,13 +408,13 @@ func TestClient_ImportRoaring(t *testing.T) { // Send import request. host := cluster[0].URL() c := MustNewClient(host, http.GetHTTPClient(nil)) - roaringData, _ := hex.DecodeString("3B3001000100000900010000000100010009000100") + roaringData, _ := hex.DecodeString("3B3001000100000900010000000100010009000100") // [1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 65537] if err := c.ImportRoaring(context.Background(), &cluster[0].API.Node().URI, "i", "f", 0, false, roaringData); err != nil { t.Fatal(err) } hldr := test.Holder{Holder: cluster[0].Server.Holder()} - // Verify data. + // Verify data on node 0. if a := hldr.Row("i", "f", 0).Columns(); !reflect.DeepEqual(a, []uint64{1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 65537}) { t.Fatalf("unexpected columns: %+v", a) } @@ -423,13 +423,79 @@ func TestClient_ImportRoaring(t *testing.T) { } hldr2 := test.Holder{Holder: cluster[1].Server.Holder()} - // Verify data. + // Verify data on node 1. if a := hldr2.Row("i", "f", 0).Columns(); !reflect.DeepEqual(a, []uint64{1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 65537}) { t.Fatalf("unexpected columns: %+v", a) } if a := hldr2.Row("i", "f", 1).Columns(); !reflect.DeepEqual(a, []uint64{0}) { t.Fatalf("unexpected columns: %+v", a) } + + // Ensure that sending a roaring import with the clear flag works as expected. + roaringDataClear, _ := hex.DecodeString("3A30000001000000010001001000000003000400") // [65539, 65540] + if err := c.ImportRoaring(context.Background(), &cluster[0].API.Node().URI, "i", "f", 0, false, roaringDataClear, pilosa.OptImportOptionsClear(true)); err != nil { + t.Fatal(err) + } + + // Verify data on node 0. + if a := hldr.Row("i", "f", 0).Columns(); !reflect.DeepEqual(a, []uint64{1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 65537}) { + t.Fatalf("unexpected clear columns: %+v", a) + } + if a := hldr.Row("i", "f", 1).Columns(); !reflect.DeepEqual(a, []uint64{0}) { + t.Fatalf("unexpected clear columns: %+v", a) + } + + // Verify data on node 1. + if a := hldr2.Row("i", "f", 0).Columns(); !reflect.DeepEqual(a, []uint64{1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 65537}) { + t.Fatalf("unexpected clear columns: %+v", a) + } + if a := hldr2.Row("i", "f", 1).Columns(); !reflect.DeepEqual(a, []uint64{0}) { + t.Fatalf("unexpected clear columns: %+v", a) + } + + // Ensure that sending a roaring import with the clear flag works as expected. + roaringDataClear, _ = hex.DecodeString("3A300000020000000000010001000100180000001C0000000400060001000300") // [4, 6, 65537, 65539] + if err := c.ImportRoaring(context.Background(), &cluster[0].API.Node().URI, "i", "f", 0, false, roaringDataClear, pilosa.OptImportOptionsClear(true)); err != nil { + t.Fatal(err) + } + + // Verify data on node 0. + if a := hldr.Row("i", "f", 0).Columns(); !reflect.DeepEqual(a, []uint64{1, 2, 3, 5, 7, 8, 9, 10}) { + t.Fatalf("unexpected clear columns: %+v", a) + } + if a := hldr.Row("i", "f", 1).Columns(); !reflect.DeepEqual(a, []uint64{0}) { + t.Fatalf("unexpected clear columns: %+v", a) + } + + // Verify data on node 1. + if a := hldr2.Row("i", "f", 0).Columns(); !reflect.DeepEqual(a, []uint64{1, 2, 3, 5, 7, 8, 9, 10}) { + t.Fatalf("unexpected clear columns: %+v", a) + } + if a := hldr2.Row("i", "f", 1).Columns(); !reflect.DeepEqual(a, []uint64{0}) { + t.Fatalf("unexpected clear columns: %+v", a) + } + + // Ensure that sending a roaring import with the clear flag works as expected. + roaringDataClear, _ = hex.DecodeString("3B3001000100000900010000000100010009000100") // [1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 65537] + if err := c.ImportRoaring(context.Background(), &cluster[0].API.Node().URI, "i", "f", 0, false, roaringDataClear, pilosa.OptImportOptionsClear(true)); err != nil { + t.Fatal(err) + } + + // Verify data on node 0. + if a := hldr.Row("i", "f", 0).Columns(); !reflect.DeepEqual(a, []uint64{}) { + t.Fatalf("unexpected clear columns: %+v", a) + } + if a := hldr.Row("i", "f", 1).Columns(); !reflect.DeepEqual(a, []uint64{0}) { + t.Fatalf("unexpected clear columns: %+v", a) + } + + // Verify data on node 1. + if a := hldr2.Row("i", "f", 0).Columns(); !reflect.DeepEqual(a, []uint64{}) { + t.Fatalf("unexpected clear columns: %+v", a) + } + if a := hldr2.Row("i", "f", 1).Columns(); !reflect.DeepEqual(a, []uint64{0}) { + t.Fatalf("unexpected clear columns: %+v", a) + } } // Ensure client can bulk import data. diff --git a/http/handler.go b/http/handler.go index 24ea6e76f..6007d72a8 100644 --- a/http/handler.go +++ b/http/handler.go @@ -1511,6 +1511,9 @@ func (h *Handler) handlePostImportRoaring(w http.ResponseWriter, r *http.Request remote = true } + // If the clear flag is true, treat the import as clear bits. + doClear := q.Get("clear") == "true" + // Read entire body. body, err := ioutil.ReadAll(r.Body) if err != nil { @@ -1527,7 +1530,7 @@ func (h *Handler) handlePostImportRoaring(w http.ResponseWriter, r *http.Request resp := &pilosa.ImportResponse{} // TODO give meaningful stats for import - err = h.api.ImportRoaring(r.Context(), urlVars["index"], urlVars["field"], shard, remote, body) + err = h.api.ImportRoaring(r.Context(), urlVars["index"], urlVars["field"], shard, remote, body, pilosa.OptImportOptionsClear(doClear)) if err != nil { resp.Err = err.Error() if _, ok := err.(pilosa.BadRequestError); ok {