From caf8e067129b8c7fdcbd44c6b1c937dcc67cbc74 Mon Sep 17 00:00:00 2001 From: Travis Turner Date: Mon, 22 Oct 2018 13:35:32 -0500 Subject: [PATCH 1/7] 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 2/7] 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 3/7] 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 4/7] 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 5/7] 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 6/7] 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 7/7] 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 {