From 12d608c80dc6a9d97978d94a68cb5bc676e49b26 Mon Sep 17 00:00:00 2001 From: Travis Turner Date: Fri, 9 Dec 2022 14:24:47 -0600 Subject: [PATCH] Move batch.Importer interface to pilosa.Importer (#2347) * Move batch.Importer interface to pilosa.Importer In addition to moving the interface, it updates all the methods to use dax.TableID (for example) intead of a string pilosa index name. * Change unused onPremImporter methods to no-op. onPremImporter is a wrapper around API which implements the Importer interface. This is currently only used by sql3 running locally in standard (i.e not "serverless") mode. Because sql3 always sets `useShardTransactionalEndpoint = true`, There are several methods which this implemtation of the Importer interface does not use, and therefore they intentionally no-op. --- batch/batch.go | 82 +++++++----- batch/batch_test.go | 84 ++++++------ batch/importer.go | 211 ------------------------------- client/api.go | 24 +++- client/importer.go | 181 +++++++++++++++++++++----- dax/queryer/importer.go | 129 ------------------- dax/queryer/queryer.go | 2 +- idk/ingest.go | 16 ++- idk/ingest_test.go | 5 +- idk/mds/importer.go | 65 +++++----- importer.go | 88 +++++++++++++ schema.go | 6 - server/server.go | 5 +- sql3/planner/executionplanner.go | 5 +- sql3/planner/opinsert.go | 2 +- 15 files changed, 401 insertions(+), 504 deletions(-) delete mode 100644 batch/importer.go delete mode 100644 dax/queryer/importer.go create mode 100644 importer.go diff --git a/batch/batch.go b/batch/batch.go index 270b21e63..b4e63f95b 100644 --- a/batch/batch.go +++ b/batch/batch.go @@ -12,6 +12,7 @@ import ( featurebase "github.com/molecula/featurebase/v3" "github.com/molecula/featurebase/v3/batch/egpool" + "github.com/molecula/featurebase/v3/dax" "github.com/molecula/featurebase/v3/logger" "github.com/molecula/featurebase/v3/pql" "github.com/molecula/featurebase/v3/roaring" @@ -95,8 +96,8 @@ type agedTranslation struct { // // nil values are ignored. type Batch struct { - importer Importer - index *featurebase.IndexInfo + importer featurebase.Importer + tbl *dax.Table header []*featurebase.FieldInfo headerMap map[string]*featurebase.FieldInfo @@ -234,7 +235,7 @@ func OptUseShardTransactionalEndpoint(use bool) BatchOption { } } -func OptImporter(i Importer) BatchOption { +func OptImporter(i featurebase.Importer) BatchOption { return func(b *Batch) error { b.importer = i return nil @@ -245,7 +246,7 @@ func OptImporter(i Importer) BatchOption { // index, set of fields, and will take "size" records before returning // ErrBatchNowFull. The positions of the Fields in 'fields' correspond to the // positions of values in the Row's Values passed to Batch.Add(). -func NewBatch(importer Importer, size int, index *featurebase.IndexInfo, fields []*featurebase.FieldInfo, opts ...BatchOption) (*Batch, error) { +func NewBatch(importer featurebase.Importer, size int, tbl *dax.Table, fields []*featurebase.FieldInfo, opts ...BatchOption) (*Batch, error) { if len(fields) == 0 { return nil, errors.New("can't batch with no fields") } else if size == 0 { @@ -301,7 +302,7 @@ func NewBatch(importer Importer, size int, index *featurebase.IndexInfo, fields header: fields, headerMap: headerMap, prevDuration: time.Minute * 11, - index: index, + tbl: tbl, ids: make([]uint64, 0, size), rowIDs: rowIDs, clearRowIDs: make(map[int]map[int]uint64), @@ -868,7 +869,7 @@ func (b *Batch) doTranslation() error { // Create the keys. start := time.Now() - trans, err := b.createIndexKeys(b.index, keys...) + trans, err := b.createIndexKeys(keys...) if err != nil { return errors.Wrap(err, "translating col keys") } @@ -1065,12 +1066,12 @@ func (b *Batch) doTranslation() error { return eg.Wait() } -func (b *Batch) createIndexKeys(index *featurebase.IndexInfo, keys ...string) (map[string]uint64, error) { +func (b *Batch) createIndexKeys(keys ...string) (map[string]uint64, error) { ctx := context.Background() batchSize := b.keyTranslateBatchSize if batchSize <= 0 || len(keys) <= batchSize { - return b.importer.CreateIndexKeys(ctx, index, keys...) + return b.importer.CreateTableKeys(ctx, b.tbl.ID, keys...) } results := make(map[string]uint64, len(keys)) @@ -1080,7 +1081,7 @@ func (b *Batch) createIndexKeys(index *featurebase.IndexInfo, keys ...string) (m keySlice = keySlice[:batchSize] } - trans, err := b.importer.CreateIndexKeys(ctx, index, keySlice...) + trans, err := b.importer.CreateTableKeys(ctx, b.tbl.ID, keySlice...) if err != nil { return nil, err } else if len(trans) != len(keySlice) { @@ -1101,7 +1102,7 @@ func (b *Batch) createFieldKeys(field *featurebase.FieldInfo, keys ...string) (m batchSize := b.keyTranslateBatchSize if batchSize <= 0 || len(keys) <= batchSize { - return b.importer.CreateFieldKeys(ctx, b.index.Name, field, keys...) + return b.importer.CreateFieldKeys(ctx, b.tbl.ID, dax.FieldName(field.Name), keys...) } results := make(map[string]uint64, len(keys)) @@ -1111,7 +1112,7 @@ func (b *Batch) createFieldKeys(field *featurebase.FieldInfo, keys ...string) (m keySlice = keySlice[:batchSize] } - trans, err := b.importer.CreateFieldKeys(ctx, b.index.Name, field, keySlice...) + trans, err := b.importer.CreateFieldKeys(ctx, b.tbl.ID, dax.FieldName(field.Name), keySlice...) if err != nil { return nil, err } else if len(trans) != len(keySlice) { @@ -1191,7 +1192,7 @@ func (b *Batch) doImportShardTransactional(frags, clearFrags fragments) error { shard := shard request := request eg.Go(func() error { - return b.importer.ImportRoaringShard(ctx, b.index.Name, shard, request) + return b.importer.ImportRoaringShard(ctx, b.tbl.ID, shard, request) }) } err := eg.Wait() @@ -1224,17 +1225,25 @@ func (b *Batch) doImport(frags, clearFrags fragments) error { clearViewMap := clearFrags.GetViewMap(shard, field) if len(clearViewMap) > 0 { startx := time.Now() - err := b.importer.ImportRoaringBitmap(ctx, b.index.Name, b.indexField(field), shard, clearViewMap, true) + fld, err := b.tableField(dax.FieldName(field)) if err != nil { + return errors.Wrapf(err, "getting tablefield: %s", field) + } + if err := b.importer.ImportRoaringBitmap(ctx, b.tbl.ID, fld, shard, clearViewMap, true); err != nil { return errors.Wrapf(err, "import clearing clearing data for %s", field) } b.log.Debugf("imp-roar-clr %s,shard:%d,views:%d %v", field, shard, len(clearViewMap), time.Since(startx)) } starty := time.Now() - err := b.importer.ImportRoaringBitmap(ctx, b.index.Name, b.indexField(field), shard, viewMap, false) + fld, err := b.tableField(dax.FieldName(field)) + if err != nil { + return errors.Wrapf(err, "getting tablefield: %s", field) + } + + ferr := b.importer.ImportRoaringBitmap(ctx, b.tbl.ID, fld, shard, viewMap, false) b.log.Debugf("imp-roar %s,shard:%d,views:%d %v", field, shard, len(clearViewMap), time.Since(starty)) - return errors.Wrapf(err, "importing data for %s", field) + return errors.Wrapf(ferr, "importing data for %s", field) }) } eg.Go(func() error { return b.importValueData() }) @@ -1251,20 +1260,26 @@ func (b *Batch) doImport(frags, clearFrags fragments) error { return nil } -// indexField is a helper function which was introduced when we switched the +// tableField is a helper function which was introduced when we switched the // index and field types from being client types (e.g client.Index, // client.Field) to being featurebase types (e.g. featurebase.IndexInfo, -// featurebase.FieldInfo). Unlike client.Index, featurebase.IndexInfo is not -// expected to contain the "_exists" field. So calling Field("_exists") on -// IndexInfo results in a nil field. This method creates an instance of -// FieldInfo for the "_exists" field. -func (b *Batch) indexField(field string) *featurebase.FieldInfo { - if field == existenceFieldName { - return &featurebase.FieldInfo{ +// featurebase.FieldInfo), and then to being dax types (e.g. dax.Table, +// dax.Field). Unlike client.Index, featurebase.IndexInfo (and also dax.Table) +// is not expected to contain the "_exists" field. So calling Field("_exists") +// on IndexInfo results in a nil field. This method creates an instance of +// dax.Field for the "_exists" field. +func (b *Batch) tableField(fname dax.FieldName) (*dax.Field, error) { + if fname == existenceFieldName { + return &dax.Field{ Name: existenceFieldName, - } + }, nil } - return b.index.Field(field) + + fld, ok := b.tbl.Field(fname) + if !ok { + return nil, errors.Errorf("field not in table: %s", fname) + } + return fld, nil } func anyCause(cause error, errs ...error) error { @@ -1299,7 +1314,10 @@ func (b *Batch) makeFragments(frags, clearFrags fragments) (fragments, fragments emptyClearRows := make(map[int]uint64) // create _exists fragments if needed - if b.index.Options.TrackExistence { + // TODO(tlt): maybe make this a separate flag for backward compatibility? + // (because dax.Table doesn't have this). + //if b.index.Options.TrackExistence { + if true { var curBM *roaring.Bitmap curShard := ^uint64(0) // impossible sentinel value for shard. for _, col := range b.ids { @@ -1678,13 +1696,15 @@ func (b *Batch) importValueData() error { endIdx := i shard := curShard field := b.headerMap[fieldName] - path, data, err := b.importer.EncodeImportValues(ctx, b.index.Name, field, shard, bvalues[startIdx:endIdx], ids[startIdx:endIdx], false) + fld := featurebase.FieldInfoToField(field) + path, data, err := b.importer.EncodeImportValues(ctx, b.tbl.ID, fld, shard, bvalues[startIdx:endIdx], ids[startIdx:endIdx], false) if err != nil { return errors.Wrap(err, "encoding import values") } eg.Go(func() error { start := time.Now() - err := b.importer.DoImport(ctx, b.index.Name, field, shard, path, data) + fld := featurebase.FieldInfoToField(field) + err := b.importer.DoImport(ctx, b.tbl.ID, fld, shard, path, data) b.log.Debugf("imp-vals %s,shard:%d,data:%d %v", field, shard, len(data), time.Since(start)) return errors.Wrapf(err, "importing values for field = %s", field.Name) }) @@ -1771,13 +1791,15 @@ func (b *Batch) importMutexData() error { endIdx := i shard := curShard field := field - path, data, err := b.importer.EncodeImport(ctx, b.index.Name, field, shard, rowIDs[startIdx:endIdx], ids[startIdx:endIdx], false) + fld := featurebase.FieldInfoToField(field) + path, data, err := b.importer.EncodeImport(ctx, b.tbl.ID, fld, shard, rowIDs[startIdx:endIdx], ids[startIdx:endIdx], false) if err != nil { return errors.Wrap(err, "encoding mutex import") } eg.Go(func() error { start := time.Now() - err := b.importer.DoImport(ctx, b.index.Name, field, shard, path, data) + fld := featurebase.FieldInfoToField(field) + err := b.importer.DoImport(ctx, b.tbl.ID, fld, shard, path, data) b.log.Debugf("imp-mux %s,shard:%d,data:%d %v", field.Name, shard, len(data), time.Since(start)) return errors.Wrapf(err, "importing values for field = %s", field.Name) }) diff --git a/batch/batch_test.go b/batch/batch_test.go index 40625dddf..025a52e78 100644 --- a/batch/batch_test.go +++ b/batch/batch_test.go @@ -24,9 +24,9 @@ func TestAgainstCluster(t *testing.T) { cli, err := client.NewClient("featurebase:10101") assert.NoError(t, err) - importer := client.NewImporter(cli) sapi := client.NewSchemaAPI(cli) qapi := client.NewQueryAPI(cli) + importer := client.NewImporter(cli, sapi) t.Run("string-slice-combos", func(t *testing.T) { testStringSliceCombos(t, importer, sapi, qapi) }) t.Run("import-batch-ints", func(t *testing.T) { testImportBatchInts(t, importer, sapi, qapi) }) @@ -49,7 +49,7 @@ func TestAgainstCluster(t *testing.T) { t.Run("test-mutex-nil-clear-key", func(t *testing.T) { mutexNilClearKey(t, importer, sapi, qapi) }) } -func testStringSliceCombos(t *testing.T, importer Importer, sapi featurebase.SchemaAPI, qapi featurebase.QueryAPI) { +func testStringSliceCombos(t *testing.T, importer featurebase.Importer, sapi featurebase.SchemaAPI, qapi featurebase.QueryAPI) { ctx := context.Background() idx := &featurebase.IndexInfo{ @@ -76,7 +76,7 @@ func testStringSliceCombos(t *testing.T, importer Importer, sapi featurebase.Sch assert.NoError(t, sapi.DeleteTable(ctx, tbl.Name)) }() - b, err := NewBatch(importer, 5, idx, idx.Fields) + b, err := NewBatch(importer, 5, tbl, idx.Fields) if err != nil { t.Fatalf("creating new batch: %v", err) } @@ -233,7 +233,7 @@ func ingestRecords(records []Row, batch *Batch) error { return nil } -func testImportBatchInts(t *testing.T, importer Importer, sapi featurebase.SchemaAPI, qapi featurebase.QueryAPI) { +func testImportBatchInts(t *testing.T, importer featurebase.Importer, sapi featurebase.SchemaAPI, qapi featurebase.QueryAPI) { ctx := context.Background() idx := &featurebase.IndexInfo{ @@ -258,7 +258,7 @@ func testImportBatchInts(t *testing.T, importer Importer, sapi featurebase.Schem assert.NoError(t, sapi.DeleteTable(ctx, tbl.Name)) }() - b, err := NewBatch(importer, 3, idx, idx.Fields) + b, err := NewBatch(importer, 3, tbl, idx.Fields) if err != nil { t.Fatalf("getting batch: %v", err) } @@ -309,7 +309,7 @@ func testImportBatchInts(t *testing.T, importer Importer, sapi featurebase.Schem } } -func testImportBatchSorting(t *testing.T, importer Importer, sapi featurebase.SchemaAPI, qapi featurebase.QueryAPI) { +func testImportBatchSorting(t *testing.T, importer featurebase.Importer, sapi featurebase.SchemaAPI, qapi featurebase.QueryAPI) { ctx := context.Background() idx := &featurebase.IndexInfo{ @@ -342,7 +342,7 @@ func testImportBatchSorting(t *testing.T, importer Importer, sapi featurebase.Sc assert.NoError(t, sapi.DeleteTable(ctx, tbl.Name)) }() - b, err := NewBatch(importer, 100, idx, idx.Fields) + b, err := NewBatch(importer, 100, tbl, idx.Fields) if err != nil { t.Fatalf("getting batch: %v", err) } @@ -382,7 +382,7 @@ func testImportBatchSorting(t *testing.T, importer Importer, sapi featurebase.Sc assert.Equal(t, uint64(100), count) } -func testTrimNull(t *testing.T, importer Importer, sapi featurebase.SchemaAPI, qapi featurebase.QueryAPI) { +func testTrimNull(t *testing.T, importer featurebase.Importer, sapi featurebase.SchemaAPI, qapi featurebase.QueryAPI) { ctx := context.Background() fieldName := "empty" @@ -408,7 +408,7 @@ func testTrimNull(t *testing.T, importer Importer, sapi featurebase.SchemaAPI, q assert.NoError(t, sapi.DeleteTable(ctx, tbl.Name)) }() - b, err := NewBatch(importer, 3, idx, idx.Fields) + b, err := NewBatch(importer, 3, tbl, idx.Fields) if err != nil { t.Fatalf("getting batch: %v", err) } @@ -440,7 +440,7 @@ func testTrimNull(t *testing.T, importer Importer, sapi featurebase.SchemaAPI, q assert.Equal(t, []uint64{}, row.Columns()) } - b, err = NewBatch(importer, 4, idx, idx.Fields) + b, err = NewBatch(importer, 4, tbl, idx.Fields) if err != nil { t.Fatalf("getting batch: %v", err) } @@ -499,7 +499,7 @@ func testTrimNull(t *testing.T, importer Importer, sapi featurebase.SchemaAPI, q } } -func testStringSliceEmptyAndNil(t *testing.T, importer Importer, sapi featurebase.SchemaAPI, qapi featurebase.QueryAPI) { +func testStringSliceEmptyAndNil(t *testing.T, importer featurebase.Importer, sapi featurebase.SchemaAPI, qapi featurebase.QueryAPI) { ctx := context.Background() idx := &featurebase.IndexInfo{ @@ -529,7 +529,7 @@ func testStringSliceEmptyAndNil(t *testing.T, importer Importer, sapi featurebas // first create a batch and test adding a single value with empty // string - this failed with a translation error at one point, and // how we catch it and treat it like a nil. - b, err := NewBatch(importer, 2, idx, idx.Fields) + b, err := NewBatch(importer, 2, tbl, idx.Fields) if err != nil { t.Fatalf("creating new batch: %v", err) } @@ -546,7 +546,7 @@ func testStringSliceEmptyAndNil(t *testing.T, importer Importer, sapi featurebas } // now create a batch and add a mixture of string slice values - b, err = NewBatch(importer, 6, idx, idx.Fields) + b, err = NewBatch(importer, 6, tbl, idx.Fields) if err != nil { t.Fatalf("creating new batch: %v", err) } @@ -625,7 +625,7 @@ func testStringSliceEmptyAndNil(t *testing.T, importer Importer, sapi featurebas } } -func testStringSlice(t *testing.T, importer Importer, sapi featurebase.SchemaAPI, qapi featurebase.QueryAPI) { +func testStringSlice(t *testing.T, importer featurebase.Importer, sapi featurebase.SchemaAPI, qapi featurebase.QueryAPI) { ctx := context.Background() idx := &featurebase.IndexInfo{ @@ -652,7 +652,7 @@ func testStringSlice(t *testing.T, importer Importer, sapi featurebase.SchemaAPI assert.NoError(t, sapi.DeleteTable(ctx, tbl.Name)) }() - b, err := NewBatch(importer, 3, idx, idx.Fields) + b, err := NewBatch(importer, 3, tbl, idx.Fields) if err != nil { t.Fatalf("creating new batch: %v", err) } @@ -750,7 +750,7 @@ func testStringSlice(t *testing.T, importer Importer, sapi featurebase.SchemaAPI assert.Equal(t, []uint64{0, 1}, row.Columns()) } -func testSingleClearBatchRegression(t *testing.T, importer Importer, sapi featurebase.SchemaAPI, qapi featurebase.QueryAPI) { +func testSingleClearBatchRegression(t *testing.T, importer featurebase.Importer, sapi featurebase.SchemaAPI, qapi featurebase.QueryAPI) { ctx := context.Background() idx := &featurebase.IndexInfo{ @@ -783,7 +783,7 @@ func testSingleClearBatchRegression(t *testing.T, importer Importer, sapi featur Query: "Set(1, zero='row1')", }) - b, err := NewBatch(importer, 1, idx, idx.Fields) + b, err := NewBatch(importer, 1, tbl, idx.Fields) if err != nil { t.Fatalf("getting new batch: %v", err) } @@ -809,7 +809,7 @@ func testSingleClearBatchRegression(t *testing.T, importer Importer, sapi featur assert.Equal(t, []uint64{}, row.Columns()) } -func testBatches(t *testing.T, importer Importer, sapi featurebase.SchemaAPI, qapi featurebase.QueryAPI) { +func testBatches(t *testing.T, importer featurebase.Importer, sapi featurebase.SchemaAPI, qapi featurebase.QueryAPI) { ctx := context.Background() idx := &featurebase.IndexInfo{ @@ -869,7 +869,7 @@ func testBatches(t *testing.T, importer Importer, sapi featurebase.SchemaAPI, qa assert.NoError(t, sapi.DeleteTable(ctx, tbl.Name)) }() - b, err := NewBatch(importer, 10, idx, idx.Fields) + b, err := NewBatch(importer, 10, tbl, idx.Fields) if err != nil { t.Fatalf("getting new batch: %v", err) } @@ -1264,7 +1264,7 @@ func testBatches(t *testing.T, importer Importer, sapi featurebase.SchemaAPI, qa // TODO test importing across multiple shards } -func testBatchesStringIDs(t *testing.T, importer Importer, sapi featurebase.SchemaAPI, qapi featurebase.QueryAPI) { +func testBatchesStringIDs(t *testing.T, importer featurebase.Importer, sapi featurebase.SchemaAPI, qapi featurebase.QueryAPI) { ctx := context.Background() idx := &featurebase.IndexInfo{ @@ -1309,7 +1309,7 @@ func testBatchesStringIDs(t *testing.T, importer Importer, sapi featurebase.Sche assert.NoError(t, sapi.DeleteTable(ctx, tbl.Name)) }() - b, err := NewBatch(importer, 3, idx, idx.Fields) + b, err := NewBatch(importer, 3, tbl, idx.Fields) if err != nil { t.Fatalf("getting new batch: %v", err) } @@ -1568,7 +1568,7 @@ func TestQuantizedTime(t *testing.T) { } } -func testBatchStaleness(t *testing.T, importer Importer, sapi featurebase.SchemaAPI, qapi featurebase.QueryAPI) { +func testBatchStaleness(t *testing.T, importer featurebase.Importer, sapi featurebase.SchemaAPI, qapi featurebase.QueryAPI) { ctx := context.Background() idx := &featurebase.IndexInfo{ @@ -1594,7 +1594,7 @@ func testBatchStaleness(t *testing.T, importer Importer, sapi featurebase.Schema assert.NoError(t, sapi.DeleteTable(ctx, tbl.Name)) }() - b, err := NewBatch(importer, 3, idx, idx.Fields, OptMaxStaleness(time.Millisecond)) + b, err := NewBatch(importer, 3, tbl, idx.Fields, OptMaxStaleness(time.Millisecond)) if err != nil { t.Fatalf("getting batch: %v", err) } @@ -1615,7 +1615,7 @@ func testBatchStaleness(t *testing.T, importer Importer, sapi featurebase.Schema } } -func testImportBatchMultipleInts(t *testing.T, importer Importer, sapi featurebase.SchemaAPI, qapi featurebase.QueryAPI) { +func testImportBatchMultipleInts(t *testing.T, importer featurebase.Importer, sapi featurebase.SchemaAPI, qapi featurebase.QueryAPI) { ctx := context.Background() idx := &featurebase.IndexInfo{ @@ -1641,7 +1641,7 @@ func testImportBatchMultipleInts(t *testing.T, importer Importer, sapi featureba assert.NoError(t, sapi.DeleteTable(ctx, tbl.Name)) }() - b, err := NewBatch(importer, 6, idx, idx.Fields, OptUseShardTransactionalEndpoint(true)) + b, err := NewBatch(importer, 6, tbl, idx.Fields, OptUseShardTransactionalEndpoint(true)) if err != nil { t.Fatalf("getting batch: %v", err) } @@ -1676,7 +1676,7 @@ func testImportBatchMultipleInts(t *testing.T, importer Importer, sapi featureba // testImportBatchMultipleTimestamps tests if nils are handles correctly for TS // in batch imports -func testImportBatchMultipleTimestamps(t *testing.T, importer Importer, sapi featurebase.SchemaAPI, qapi featurebase.QueryAPI) { +func testImportBatchMultipleTimestamps(t *testing.T, importer featurebase.Importer, sapi featurebase.SchemaAPI, qapi featurebase.QueryAPI) { ctx := context.Background() fieldName := "ts2" @@ -1702,11 +1702,11 @@ func testImportBatchMultipleTimestamps(t *testing.T, importer Importer, sapi fea assert.NoError(t, sapi.DeleteTable(ctx, tbl.Name)) }() - b1, err := NewBatch(importer, 6, idx, idx.Fields) + b1, err := NewBatch(importer, 6, tbl, idx.Fields) if err != nil { t.Fatalf("getting batch 1: %v", err) } - b2, err := NewBatch(importer, 6, idx, idx.Fields, OptUseShardTransactionalEndpoint(true)) + b2, err := NewBatch(importer, 6, tbl, idx.Fields, OptUseShardTransactionalEndpoint(true)) if err != nil { t.Fatalf("getting batch 2: %v", err) } @@ -1770,7 +1770,7 @@ func testImportBatchMultipleTimestamps(t *testing.T, importer Importer, sapi fea } } -func testImportBatchSetsAndClears(t *testing.T, importer Importer, sapi featurebase.SchemaAPI, qapi featurebase.QueryAPI) { +func testImportBatchSetsAndClears(t *testing.T, importer featurebase.Importer, sapi featurebase.SchemaAPI, qapi featurebase.QueryAPI) { ctx := context.Background() idx := &featurebase.IndexInfo{ @@ -1796,7 +1796,7 @@ func testImportBatchSetsAndClears(t *testing.T, importer Importer, sapi featureb assert.NoError(t, sapi.DeleteTable(ctx, tbl.Name)) }() - b, err := NewBatch(importer, 6, idx, idx.Fields, OptUseShardTransactionalEndpoint(true)) + b, err := NewBatch(importer, 6, tbl, idx.Fields, OptUseShardTransactionalEndpoint(true)) if err != nil { t.Fatalf("getting batch: %v", err) } @@ -1860,7 +1860,7 @@ func testImportBatchSetsAndClears(t *testing.T, importer Importer, sapi featureb // it didn't get removed from the cache because a full recalculation // had no way to clear the cache, it would just reset existing // values. We added Clear on the cache interface to fix this. -func testTopNCacheRegression(t *testing.T, importer Importer, sapi featurebase.SchemaAPI, qapi featurebase.QueryAPI) { +func testTopNCacheRegression(t *testing.T, importer featurebase.Importer, sapi featurebase.SchemaAPI, qapi featurebase.QueryAPI) { ctx := context.Background() idx := &featurebase.IndexInfo{ @@ -1886,7 +1886,7 @@ func testTopNCacheRegression(t *testing.T, importer Importer, sapi featurebase.S assert.NoError(t, sapi.DeleteTable(ctx, tbl.Name)) }() - b, err := NewBatch(importer, 3, idx, idx.Fields, OptUseShardTransactionalEndpoint(true)) + b, err := NewBatch(importer, 3, tbl, idx.Fields, OptUseShardTransactionalEndpoint(true)) if err != nil { t.Fatalf("getting batch: %v", err) } @@ -1944,7 +1944,7 @@ func testTopNCacheRegression(t *testing.T, importer Importer, sapi featurebase.S // with different values that only the last value is set and the bits aren't // mixed together. It adds a different ID in between the two same ones which // triggered a bug because we were sorting by shard rather than ID. -func testMultipleIntSameBatch(t *testing.T, importer Importer, sapi featurebase.SchemaAPI, qapi featurebase.QueryAPI) { +func testMultipleIntSameBatch(t *testing.T, importer featurebase.Importer, sapi featurebase.SchemaAPI, qapi featurebase.QueryAPI) { ctx := context.Background() idx := &featurebase.IndexInfo{ @@ -1969,7 +1969,7 @@ func testMultipleIntSameBatch(t *testing.T, importer Importer, sapi featurebase. assert.NoError(t, sapi.DeleteTable(ctx, tbl.Name)) }() - b, err := NewBatch(importer, 4, idx, idx.Fields, OptUseShardTransactionalEndpoint(true)) + b, err := NewBatch(importer, 4, tbl, idx.Fields, OptUseShardTransactionalEndpoint(true)) if err != nil { t.Fatalf("getting batch: %v", err) } @@ -2013,7 +2013,7 @@ func testMultipleIntSameBatch(t *testing.T, importer Importer, sapi featurebase. // one in a batch did not get any bits set in their clear bitmap, and // in fact, all the bits were set in the clear bitmap for the first // shard. -func mutexClearRegression(t *testing.T, importer Importer, sapi featurebase.SchemaAPI, qapi featurebase.QueryAPI) { +func mutexClearRegression(t *testing.T, importer featurebase.Importer, sapi featurebase.SchemaAPI, qapi featurebase.QueryAPI) { ctx := context.Background() idx := &featurebase.IndexInfo{ @@ -2039,7 +2039,7 @@ func mutexClearRegression(t *testing.T, importer Importer, sapi featurebase.Sche assert.NoError(t, sapi.DeleteTable(ctx, tbl.Name)) }() - b, err := NewBatch(importer, 11, idx, idx.Fields, OptUseShardTransactionalEndpoint(true)) + b, err := NewBatch(importer, 11, tbl, idx.Fields, OptUseShardTransactionalEndpoint(true)) if err != nil { t.Fatalf("getting batch: %v", err) } @@ -2094,7 +2094,7 @@ func mutexClearRegression(t *testing.T, importer Importer, sapi featurebase.Sche } // test clearing record with explict nil -func mutexNilClearID(t *testing.T, importer Importer, sapi featurebase.SchemaAPI, qapi featurebase.QueryAPI) { +func mutexNilClearID(t *testing.T, importer featurebase.Importer, sapi featurebase.SchemaAPI, qapi featurebase.QueryAPI) { ctx := context.Background() idx := &featurebase.IndexInfo{ @@ -2120,7 +2120,7 @@ func mutexNilClearID(t *testing.T, importer Importer, sapi featurebase.SchemaAPI assert.NoError(t, sapi.DeleteTable(ctx, tbl.Name)) }() - b, err := NewBatch(importer, 11, idx, idx.Fields, OptUseShardTransactionalEndpoint(true)) + b, err := NewBatch(importer, 11, tbl, idx.Fields, OptUseShardTransactionalEndpoint(true)) if err != nil { t.Fatalf("getting batch: %v", err) } @@ -2182,7 +2182,7 @@ func mutexNilClearID(t *testing.T, importer Importer, sapi featurebase.SchemaAPI } // similar test to above but with string keys -func mutexNilClearKey(t *testing.T, importer Importer, sapi featurebase.SchemaAPI, qapi featurebase.QueryAPI) { +func mutexNilClearKey(t *testing.T, importer featurebase.Importer, sapi featurebase.SchemaAPI, qapi featurebase.QueryAPI) { ctx := context.Background() idx := &featurebase.IndexInfo{ @@ -2210,7 +2210,7 @@ func mutexNilClearKey(t *testing.T, importer Importer, sapi featurebase.SchemaAP assert.NoError(t, sapi.DeleteTable(ctx, tbl.Name)) }() - b, err := NewBatch(importer, 3, idx, idx.Fields) + b, err := NewBatch(importer, 3, tbl, idx.Fields) if err != nil { t.Fatalf("getting new batch: %v", err) } @@ -2284,7 +2284,7 @@ func mutexNilClearKey(t *testing.T, importer Importer, sapi featurebase.SchemaAP assert.Equal(t, []string{"0"}, fRow.Keys) } -func testImportBatchBools(t *testing.T, importer Importer, sapi featurebase.SchemaAPI, qapi featurebase.QueryAPI) { +func testImportBatchBools(t *testing.T, importer featurebase.Importer, sapi featurebase.SchemaAPI, qapi featurebase.QueryAPI) { ctx := context.Background() idx := &featurebase.IndexInfo{ @@ -2308,7 +2308,7 @@ func testImportBatchBools(t *testing.T, importer Importer, sapi featurebase.Sche assert.NoError(t, sapi.DeleteTable(ctx, tbl.Name)) }() - b, err := NewBatch(importer, 3, idx, idx.Fields, OptUseShardTransactionalEndpoint(true)) + b, err := NewBatch(importer, 3, tbl, idx.Fields, OptUseShardTransactionalEndpoint(true)) if err != nil { t.Fatalf("getting new batch: %v", err) } diff --git a/batch/importer.go b/batch/importer.go deleted file mode 100644 index 7f1e54478..000000000 --- a/batch/importer.go +++ /dev/null @@ -1,211 +0,0 @@ -package batch - -import ( - "context" - "time" - - "github.com/golang/protobuf/proto" //nolint:staticcheck - featurebase "github.com/molecula/featurebase/v3" - featurebaseproto "github.com/molecula/featurebase/v3/encoding/proto" - "github.com/molecula/featurebase/v3/pb" - "github.com/molecula/featurebase/v3/roaring" - "github.com/pkg/errors" -) - -type Importer interface { - StartTransaction(ctx context.Context, id string, timeout time.Duration, exclusive bool, requestTimeout time.Duration) (*featurebase.Transaction, error) - FinishTransaction(ctx context.Context, id string) (*featurebase.Transaction, error) - CreateIndexKeys(ctx context.Context, idx *featurebase.IndexInfo, keys ...string) (map[string]uint64, error) - CreateFieldKeys(ctx context.Context, index string, field *featurebase.FieldInfo, keys ...string) (map[string]uint64, error) - ImportRoaringBitmap(ctx context.Context, index string, field *featurebase.FieldInfo, shard uint64, views map[string]*roaring.Bitmap, clear bool) error - ImportRoaringShard(ctx context.Context, index string, shard uint64, request *featurebase.ImportRoaringShardRequest) error - EncodeImportValues(ctx context.Context, index string, field *featurebase.FieldInfo, shard uint64, vals []int64, ids []uint64, clear bool) (path string, data []byte, err error) - EncodeImport(ctx context.Context, index string, field *featurebase.FieldInfo, shard uint64, vals, ids []uint64, clear bool) (path string, data []byte, err error) - DoImport(ctx context.Context, index string, field *featurebase.FieldInfo, shard uint64, path string, data []byte) error - - StatsTiming(name string, value time.Duration, rate float64) -} - -// Ensure type implements interface. -var _ Importer = &nopImporter{} - -// NopImporter is an implementation of the Importer interface that doesn't do -// anything. -var NopImporter Importer = &nopImporter{} - -type nopImporter struct{} - -func (n *nopImporter) StartTransaction(ctx context.Context, id string, timeout time.Duration, exclusive bool, requestTimeout time.Duration) (*featurebase.Transaction, error) { - return nil, nil -} -func (n *nopImporter) FinishTransaction(ctx context.Context, id string) (*featurebase.Transaction, error) { - return nil, nil -} -func (n *nopImporter) CreateIndexKeys(ctx context.Context, idx *featurebase.IndexInfo, keys ...string) (map[string]uint64, error) { - return nil, nil -} -func (n *nopImporter) CreateFieldKeys(ctx context.Context, index string, field *featurebase.FieldInfo, keys ...string) (map[string]uint64, error) { - return nil, nil -} -func (n *nopImporter) ImportRoaringBitmap(ctx context.Context, index string, field *featurebase.FieldInfo, shard uint64, views map[string]*roaring.Bitmap, clear bool) error { - return nil -} -func (n *nopImporter) ImportRoaringShard(ctx context.Context, index string, shard uint64, request *featurebase.ImportRoaringShardRequest) error { - return nil -} -func (n *nopImporter) EncodeImportValues(ctx context.Context, index string, field *featurebase.FieldInfo, shard uint64, vals []int64, ids []uint64, clear bool) (path string, data []byte, err error) { - return "", nil, nil -} -func (n *nopImporter) EncodeImport(ctx context.Context, index string, field *featurebase.FieldInfo, shard uint64, vals, ids []uint64, clear bool) (path string, data []byte, err error) { - return "", nil, nil -} -func (n *nopImporter) DoImport(ctx context.Context, index string, field *featurebase.FieldInfo, shard uint64, path string, data []byte) error { - return nil -} - -func (n *nopImporter) StatsTiming(name string, value time.Duration, rate float64) {} - -// Ensure type implements interface. -var _ Importer = &FeaturebaseImporter{} - -// FeaturebaseImporter is a wrapper around featurebase.API, making it a -// batch.Importer. -type FeaturebaseImporter struct { - *featurebase.API -} - -func (f *FeaturebaseImporter) StartTransaction(ctx context.Context, id string, timeout time.Duration, exclusive bool, requestTimeout time.Duration) (*featurebase.Transaction, error) { - return f.API.StartTransaction(ctx, id, timeout, exclusive, false) -} - -func (f *FeaturebaseImporter) FinishTransaction(ctx context.Context, id string) (*featurebase.Transaction, error) { - return f.API.FinishTransaction(ctx, id, false) -} - -func (f *FeaturebaseImporter) CreateIndexKeys(ctx context.Context, idx *featurebase.IndexInfo, keys ...string) (map[string]uint64, error) { - return f.API.CreateIndexKeys(ctx, idx.Name, keys...) -} - -func (f *FeaturebaseImporter) CreateFieldKeys(ctx context.Context, index string, field *featurebase.FieldInfo, keys ...string) (map[string]uint64, error) { - return f.API.CreateFieldKeys(ctx, index, field.Name, keys...) -} - -func (f *FeaturebaseImporter) ImportRoaringBitmap(ctx context.Context, index string, field *featurebase.FieldInfo, shard uint64, views map[string]*roaring.Bitmap, clear bool) error { - vs := make(map[string][]byte) - for k, v := range views { - data := roaring.BitmapsToRoaring([]*roaring.Bitmap{v}) - if len(data) > 0 { - vs[k] = data - } - } - req := &featurebase.ImportRoaringRequest{ - IndexCreatedAt: 0, - FieldCreatedAt: field.CreatedAt, - Clear: clear, - Views: vs, - } - return f.API.ImportRoaring(ctx, index, field.Name, shard, false, req) -} - -// ImportRoaringShard doesn't technically need to be implemented here, because -// the method on f.API has the same signature and already satisfies the -// interface. But we put this here to avoid possible confusion. -func (f *FeaturebaseImporter) ImportRoaringShard(ctx context.Context, index string, shard uint64, request *featurebase.ImportRoaringShardRequest) error { - return f.API.ImportRoaringShard(ctx, index, shard, request) -} - -// EncodeImportValues is kind of weird. We're trying to mimic what the client -// does here (because the Importer interface was originally based off of the -// client methods). So we end up generating a protobuf-encode byte slice. And we -// don't really use path. -func (f *FeaturebaseImporter) EncodeImportValues(ctx context.Context, index string, field *featurebase.FieldInfo, shard uint64, vals []int64, ids []uint64, clear bool) (path string, data []byte, err error) { - msg := &pb.ImportValueRequest{ - Index: index, - IndexCreatedAt: 0, - Field: field.Name, - FieldCreatedAt: field.CreatedAt, - Shard: shard, - ColumnIDs: ids, - Values: vals, - } - data, err = proto.Marshal(msg) - if err != nil { - return "", nil, errors.Wrap(err, "marshaling ImportValue to protobuf") - } - return "", data, nil -} - -func (f *FeaturebaseImporter) EncodeImport(ctx context.Context, index string, field *featurebase.FieldInfo, shard uint64, vals, ids []uint64, clear bool) (path string, data []byte, err error) { - msg := &pb.ImportRequest{ - Index: index, - IndexCreatedAt: 0, - Field: field.Name, - FieldCreatedAt: field.CreatedAt, - Shard: shard, - RowIDs: vals, - ColumnIDs: ids, - } - data, err = proto.Marshal(msg) - if err != nil { - return "", nil, errors.Wrap(err, "marshaling Import to protobuf") - } - return "", data, nil -} - -func (f *FeaturebaseImporter) DoImport(ctx context.Context, index string, field *featurebase.FieldInfo, shard uint64, path string, data []byte) error { - serializer := featurebaseproto.Serializer{} - - // Unmarshal request based on field type. - switch field.Options.Type { - case featurebase.FieldTypeInt, featurebase.FieldTypeDecimal, featurebase.FieldTypeTimestamp: - // Marshal into request object. - req := &featurebase.ImportValueRequest{} - if err := serializer.Unmarshal(data, req); err != nil { - return errors.Wrap(err, "unmarshaling import value request") - } - - qcx := f.API.Txf().NewQcx() - defer qcx.Abort() - - opts := []featurebase.ImportOption{ - featurebase.OptImportOptionsClear(req.Clear), - } - - if err := f.API.ImportValue(ctx, qcx, req, opts...); err != nil { - return errors.Wrap(err, "importing import value request") - } - - if err := qcx.Finish(); err != nil { - return errors.Wrap(err, "finishing qcx") - } - - default: - // Marshal into request object. - req := &featurebase.ImportRequest{} - if err := serializer.Unmarshal(data, req); err != nil { - return errors.Wrap(err, "unmarshaling import request") - } - - qcx := f.API.Txf().NewQcx() - defer qcx.Abort() - - opts := []featurebase.ImportOption{ - featurebase.OptImportOptionsClear(req.Clear), - } - if len(req.RowIDs) > 0 { - opts = append(opts, featurebase.OptImportOptionsIgnoreKeyCheck(true)) - } - - if err := f.API.Import(ctx, qcx, req, opts...); err != nil { - return errors.Wrap(err, "importing import request") - } - - if err := qcx.Finish(); err != nil { - return errors.Wrap(err, "finishing qcx") - } - } - - return nil -} - -func (f *FeaturebaseImporter) StatsTiming(name string, value time.Duration, rate float64) {} diff --git a/client/api.go b/client/api.go index e46bf2635..4a7f0cddc 100644 --- a/client/api.go +++ b/client/api.go @@ -27,13 +27,27 @@ func NewSchemaAPI(c *Client) *schemaAPI { } func (s *schemaAPI) TableByName(ctx context.Context, tname dax.TableName) (*dax.Table, error) { - return nil, nil + return nil, errors.New(errors.ErrUncoded, "schemaAPI.TableByName not implemented") } func (s *schemaAPI) TableByID(ctx context.Context, tid dax.TableID) (*dax.Table, error) { - return nil, nil + schema, err := s.client.Schema() + if err != nil { + return nil, errors.Wrap(err, "getting schema") + } + + idx := schema.Index(string(tid)) + if idx == nil { + return nil, errors.Errorf("index not found: %s", tid) + } + + ii := FromClientIndex(idx) + tbl := featurebase.IndexInfoToTable(ii) + + return tbl, nil } + func (s *schemaAPI) Tables(ctx context.Context) ([]*dax.Table, error) { - return nil, nil + return nil, errors.New(errors.ErrUncoded, "schemaAPI.Tables not implemented") } func (s *schemaAPI) CreateTable(ctx context.Context, tbl *dax.Table) error { @@ -67,14 +81,14 @@ func (s *schemaAPI) CreateTable(ctx context.Context, tbl *dax.Table) error { return nil } func (s *schemaAPI) CreateField(ctx context.Context, tname dax.TableName, fld *dax.Field) error { - return nil + return errors.New(errors.ErrUncoded, "schemaAPI.CreateField not implemented") } func (s *schemaAPI) DeleteTable(ctx context.Context, tname dax.TableName) error { return s.client.DeleteIndexByName(string(tname)) } func (s *schemaAPI) DeleteField(ctx context.Context, tname dax.TableName, fname dax.FieldName) error { - return nil + return errors.New(errors.ErrUncoded, "schemaAPI.DeleteField not implemented") } func (s *schemaAPI) addFieldToIndex(idx *Index, fieldName string, ffos featurebase.FieldOptions) (*Field, error) { diff --git a/client/importer.go b/client/importer.go index 9216f9025..c42c74c4d 100644 --- a/client/importer.go +++ b/client/importer.go @@ -5,6 +5,7 @@ import ( "time" featurebase "github.com/molecula/featurebase/v3" + "github.com/molecula/featurebase/v3/dax" "github.com/molecula/featurebase/v3/roaring" "github.com/pkg/errors" ) @@ -48,6 +49,24 @@ func ToClientIndex(fi *featurebase.IndexInfo) *Index { ) } +// QTableToClientIndex +func QTableToClientIndex(qtbl *dax.QualifiedTable) *Index { + sch := NewSchema() + return sch.Index(string(qtbl.Key()), + OptIndexKeys(qtbl.StringKeys()), + OptIndexTrackExistence(true), + ) +} + +// TableToClientIndex +func TableToClientIndex(tbl *dax.Table) *Index { + sch := NewSchema() + return sch.Index(string(tbl.ID), + OptIndexKeys(tbl.StringKeys()), + OptIndexTrackExistence(true), + ) +} + // fromClientField func fromClientField(cf *Field) *featurebase.FieldInfo { return &featurebase.FieldInfo{ @@ -123,6 +142,81 @@ func ToClientField(index string, ff *featurebase.FieldInfo) (*Field, error) { return idx.Field(ff.Name, opts...), nil } +const ( + existenceFieldName = "_exists" +) + +// tableField is a helper function that returns the _exists field as a +// dax.Field when necessary (because dax.Table doesn't normally track the +// _exists field. +func tableField(tbl *dax.Table, fname dax.FieldName) (*dax.Field, error) { + if fname == existenceFieldName { + return &dax.Field{ + Name: existenceFieldName, + }, nil + } + + fld, ok := tbl.Field(fname) + if !ok { + return nil, errors.Errorf("field not in table: %s", fname) + } + return fld, nil +} + +// TableFieldToClientField +func TableFieldToClientField(qtbl *dax.QualifiedTable, fname dax.FieldName) (*Field, error) { + fld, err := tableField(&qtbl.Table, fname) + if err != nil { + return nil, errors.Wrapf(err, "getting field from table: %s", fname) + } + + sch := NewSchema() + idx := sch.Index(string(qtbl.Key())) + + opts := []FieldOption{} + + switch fld.Type { + case dax.BaseTypeBool: + opts = append(opts, + OptFieldTypeBool(), + ) + case dax.BaseTypeDecimal: + opts = append(opts, + OptFieldTypeDecimal(fld.Options.Scale, fld.Options.Min, fld.Options.Max), + ) + case dax.BaseTypeID: + opts = append(opts, + OptFieldTypeMutex(CacheType(fld.Options.CacheType), int(fld.Options.CacheSize)), + OptFieldKeys(false), + ) + case dax.BaseTypeString: + opts = append(opts, + OptFieldTypeMutex(CacheType(fld.Options.CacheType), int(fld.Options.CacheSize)), + OptFieldKeys(true), + ) + case dax.BaseTypeIDSet: + opts = append(opts, + OptFieldTypeSet(CacheType(fld.Options.CacheType), int(fld.Options.CacheSize)), + OptFieldKeys(false), + ) + case dax.BaseTypeStringSet: + opts = append(opts, + OptFieldTypeSet(CacheType(fld.Options.CacheType), int(fld.Options.CacheSize)), + OptFieldKeys(true), + ) + case dax.BaseTypeInt: + opts = append(opts, + OptFieldTypeInt(fld.Options.Min.ToInt64(0), fld.Options.Max.ToInt64(0)), + ) + case dax.BaseTypeTimestamp: + opts = append(opts, + OptFieldTypeTimestamp(fld.Options.Epoch, fld.Options.TimeUnit), + ) + } + + return idx.Field(string(fld.Name), opts...), nil +} + // FromClientFields converts a slice of client Fields to a slice of featurebase // FieldInfo. func FromClientFields(cf []*Field) []*featurebase.FieldInfo { @@ -142,78 +236,99 @@ func fromClientFieldsMap(cf map[string]*Field) []*featurebase.FieldInfo { return ff } -// We can't import the batch package into client because it results in an import -// loop. That's probably an indication that this interface implementation should -// be moved somewhere else; for example, into a sub-package of the batch package -// (since it's an implementation of one of batch's interfaces). -// var _ batch.Importer = &importer{} +/////////////////////////////////////////////////////////////// -// importer is a pilosa client which implements the batch.Importer interface. -// This wrapper is necessary because of the call into client.Stats.Timing(), and -// because the client takes client specific types (like Index and Field), but -// the interface takes FeatureBase specific types. +var _ featurebase.Importer = &importer{} + +// importer is a pilosa client which implements the featurebase.Importer +// interface. This wrapper is necessary because of the call into +// client.Stats.Timing(), and because the client takes client specific types +// (like Index and Field), but the interface takes FeatureBase (and dax) +// specific types. type importer struct { - *Client + client *Client + schemaAPI featurebase.SchemaAPI } -func NewImporter(c *Client) *importer { +func NewImporter(c *Client, sapi featurebase.SchemaAPI) *importer { return &importer{ - Client: c, + client: c, + schemaAPI: sapi, } } func (i *importer) StartTransaction(ctx context.Context, id string, timeout time.Duration, exclusive bool, requestTimeout time.Duration) (*featurebase.Transaction, error) { - return i.Client.StartTransaction(id, timeout, exclusive, requestTimeout) + return i.client.StartTransaction(id, timeout, exclusive, requestTimeout) } func (i *importer) FinishTransaction(ctx context.Context, id string) (*featurebase.Transaction, error) { - return i.Client.FinishTransaction(id) + return i.client.FinishTransaction(id) } -func (i *importer) CreateIndexKeys(ctx context.Context, idx *featurebase.IndexInfo, keys ...string) (map[string]uint64, error) { - return i.Client.CreateIndexKeys(ToClientIndex(idx), keys...) +func (i *importer) CreateTableKeys(ctx context.Context, tid dax.TableID, keys ...string) (map[string]uint64, error) { + tbl, err := i.schemaAPI.TableByID(ctx, tid) + if err != nil { + return nil, err + } + return i.client.CreateIndexKeys(TableToClientIndex(tbl), keys...) } -func (i *importer) CreateFieldKeys(ctx context.Context, index string, field *featurebase.FieldInfo, keys ...string) (map[string]uint64, error) { - fld, err := ToClientField(index, field) +func (i *importer) CreateFieldKeys(ctx context.Context, tid dax.TableID, fname dax.FieldName, keys ...string) (map[string]uint64, error) { + tbl, err := i.schemaAPI.TableByID(ctx, tid) + if err != nil { + return nil, err + } + fld, ok := tbl.Field(fname) + if !ok { + return nil, errors.Errorf("field not in table: %s", fname) + } + fi := featurebase.FieldToFieldInfo(fld) + + cfld, err := ToClientField(string(tbl.ID), fi) if err != nil { return nil, errors.Wrap(err, "converting to client field") } - return i.Client.CreateFieldKeys(fld, keys...) + return i.client.CreateFieldKeys(cfld, keys...) } -func (i *importer) ImportRoaringBitmap(ctx context.Context, index string, field *featurebase.FieldInfo, shard uint64, views map[string]*roaring.Bitmap, clear bool) error { - fld, err := ToClientField(index, field) +func (i *importer) ImportRoaringBitmap(ctx context.Context, tid dax.TableID, fld *dax.Field, shard uint64, views map[string]*roaring.Bitmap, clear bool) error { + fi := featurebase.FieldToFieldInfo(fld) + + cfld, err := ToClientField(string(tid), fi) if err != nil { return errors.Wrap(err, "converting to client field") } - return i.Client.ImportRoaringBitmap(fld, shard, views, clear) + return i.client.ImportRoaringBitmap(cfld, shard, views, clear) } -func (i *importer) ImportRoaringShard(ctx context.Context, index string, shard uint64, request *featurebase.ImportRoaringShardRequest) error { - return i.Client.ImportRoaringShard(index, shard, request) +func (i *importer) ImportRoaringShard(ctx context.Context, tid dax.TableID, shard uint64, request *featurebase.ImportRoaringShardRequest) error { + return i.client.ImportRoaringShard(string(tid), shard, request) } -func (i *importer) EncodeImportValues(ctx context.Context, index string, field *featurebase.FieldInfo, shard uint64, vals []int64, ids []uint64, clear bool) (path string, data []byte, err error) { - fld, err := ToClientField(index, field) +func (i *importer) EncodeImportValues(ctx context.Context, tid dax.TableID, fld *dax.Field, shard uint64, vals []int64, ids []uint64, clear bool) (path string, data []byte, err error) { + fi := featurebase.FieldToFieldInfo(fld) + + cfld, err := ToClientField(string(tid), fi) if err != nil { return "", nil, errors.Wrap(err, "converting to client field") } - return i.Client.EncodeImportValues(fld, shard, vals, ids, clear) + return i.client.EncodeImportValues(cfld, shard, vals, ids, clear) } -func (i *importer) EncodeImport(ctx context.Context, index string, field *featurebase.FieldInfo, shard uint64, vals, ids []uint64, clear bool) (path string, data []byte, err error) { - fld, err := ToClientField(index, field) +func (i *importer) EncodeImport(ctx context.Context, tid dax.TableID, fld *dax.Field, shard uint64, vals, ids []uint64, clear bool) (path string, data []byte, err error) { + fi := featurebase.FieldToFieldInfo(fld) + + cfld, err := ToClientField(string(tid), fi) if err != nil { return "", nil, errors.Wrap(err, "converting to client field") } - return i.Client.EncodeImport(fld, shard, vals, ids, clear) + return i.client.EncodeImport(cfld, shard, vals, ids, clear) } -func (i *importer) DoImport(ctx context.Context, index string, field *featurebase.FieldInfo, shard uint64, path string, data []byte) error { - return i.Client.DoImport(index, shard, path, data) +func (i *importer) DoImport(ctx context.Context, tid dax.TableID, fld *dax.Field, shard uint64, path string, data []byte) error { + return i.client.DoImport(string(tid), shard, path, data) } func (i *importer) StatsTiming(name string, value time.Duration, rate float64) { - i.Client.Stats.Timing(name, value, rate) + i.client.Stats.Timing(name, value, rate) } diff --git a/dax/queryer/importer.go b/dax/queryer/importer.go deleted file mode 100644 index cb684fb5b..000000000 --- a/dax/queryer/importer.go +++ /dev/null @@ -1,129 +0,0 @@ -package queryer - -import ( - "context" - "strings" - "time" - - featurebase "github.com/molecula/featurebase/v3" - "github.com/molecula/featurebase/v3/batch" - "github.com/molecula/featurebase/v3/dax" - "github.com/molecula/featurebase/v3/dax/mds/schemar" - "github.com/molecula/featurebase/v3/errors" - "github.com/molecula/featurebase/v3/roaring" -) - -// Ensure type implements interface. -var _ batch.Importer = &batchImporter{} - -func newBatchImporter(importer batch.Importer, qual dax.TableQualifier, schemar schemar.Schemar) *batchImporter { - return &batchImporter{ - importer: importer, - qual: qual, - schemar: schemar, - } -} - -// batchImporter is an implementation of the batch.Importer. It is a wrapper -// around idk/mds/Importer that can take index values which are either indexName -// (like "foo") or TableKey (like "tbl__acme__db1__foo123"). This wrapper looks -// at the value to determine if it is a TableKey or not and converts it -// appropriately. It's kind of annoying; we really need to be certain where -// we're expecting indexName vs TableKey. -type batchImporter struct { - importer batch.Importer - qual dax.TableQualifier - schemar schemar.Schemar -} - -func (b *batchImporter) StartTransaction(ctx context.Context, id string, timeout time.Duration, exclusive bool, requestTimeout time.Duration) (*featurebase.Transaction, error) { - return b.importer.StartTransaction(ctx, id, timeout, exclusive, requestTimeout) -} - -func (b *batchImporter) FinishTransaction(ctx context.Context, id string) (*featurebase.Transaction, error) { - return b.importer.FinishTransaction(ctx, id) -} - -func (b *batchImporter) CreateIndexKeys(ctx context.Context, idx *featurebase.IndexInfo, keys ...string) (map[string]uint64, error) { - // Used as an example: - // qual: [acme:db1] - // INSERT INTO foo VALUE (1, 10) - // Currently, the table name coming through IndexInfo right now is the sql - // table name (i.e. foo, not tbl__acme__db1__foo123). For SELECT queries, - // we're currently doing that conversion in the orchestrator.Execute() - // method (which means that we currently only support a single index in SQL - // queries). Therefore, we need to convert idx.Name to a TableKey. - tkey, err := b.indexToQualifiedTableKey(ctx, idx.Name) - if err != nil { - return nil, errors.Wrapf(err, "converting index to qualified table key: %s", idx.Name) - } - idx.Name = string(tkey) - - return b.importer.CreateIndexKeys(ctx, idx, keys...) -} - -func (b *batchImporter) CreateFieldKeys(ctx context.Context, index string, field *featurebase.FieldInfo, keys ...string) (map[string]uint64, error) { - tkey, err := b.indexToQualifiedTableKey(ctx, index) - if err != nil { - return nil, errors.Wrapf(err, "converting index to qualified table key: %s", index) - } - return b.importer.CreateFieldKeys(ctx, string(tkey), field, keys...) -} - -func (b *batchImporter) ImportRoaringBitmap(ctx context.Context, index string, field *featurebase.FieldInfo, shard uint64, views map[string]*roaring.Bitmap, clear bool) error { - tkey, err := b.indexToQualifiedTableKey(ctx, index) - if err != nil { - return errors.Wrapf(err, "converting index to qualified table key: %s", index) - } - return b.importer.ImportRoaringBitmap(ctx, string(tkey), field, shard, views, clear) -} - -func (b *batchImporter) ImportRoaringShard(ctx context.Context, index string, shard uint64, request *featurebase.ImportRoaringShardRequest) error { - tkey, err := b.indexToQualifiedTableKey(ctx, index) - if err != nil { - return errors.Wrapf(err, "converting index to qualified table key: %s", index) - } - return b.importer.ImportRoaringShard(ctx, string(tkey), shard, request) -} - -func (b *batchImporter) EncodeImportValues(ctx context.Context, index string, field *featurebase.FieldInfo, shard uint64, vals []int64, ids []uint64, clear bool) (path string, data []byte, err error) { - tkey, err := b.indexToQualifiedTableKey(ctx, index) - if err != nil { - return "", nil, errors.Wrapf(err, "converting index to qualified table key: %s", index) - } - return b.importer.EncodeImportValues(ctx, string(tkey), field, shard, vals, ids, clear) -} - -func (b *batchImporter) EncodeImport(ctx context.Context, index string, field *featurebase.FieldInfo, shard uint64, vals, ids []uint64, clear bool) (path string, data []byte, err error) { - tkey, err := b.indexToQualifiedTableKey(ctx, index) - if err != nil { - return "", nil, errors.Wrapf(err, "converting index to qualified table key: %s", index) - } - return b.importer.EncodeImport(ctx, string(tkey), field, shard, vals, ids, clear) -} - -func (b *batchImporter) DoImport(ctx context.Context, index string, field *featurebase.FieldInfo, shard uint64, path string, data []byte) error { - tkey, err := b.indexToQualifiedTableKey(ctx, index) - if err != nil { - return errors.Wrapf(err, "converting index to qualified table key: %s", index) - } - return b.importer.DoImport(ctx, string(tkey), field, shard, path, data) -} - -func (b *batchImporter) StatsTiming(name string, value time.Duration, rate float64) { - b.importer.StatsTiming(name, value, rate) -} - -// TODO(tlt): this method was copied from orchestrator.go. Can we centralize -// this logic? -func (b *batchImporter) indexToQualifiedTableKey(ctx context.Context, index string) (dax.TableKey, error) { - if strings.HasPrefix(index, dax.PrefixTable+dax.TableKeyDelimiter) { - return dax.TableKey(index), nil - } - - qtid, err := b.schemar.TableID(ctx, b.qual, dax.TableName(index)) - if err != nil { - return "", errors.Wrapf(err, "converting index to qualified table id: %s", index) - } - return qtid.Key(), nil -} diff --git a/dax/queryer/queryer.go b/dax/queryer/queryer.go index 7b0a98b35..667172909 100644 --- a/dax/queryer/queryer.go +++ b/dax/queryer/queryer.go @@ -129,7 +129,7 @@ func (q *Queryer) QuerySQL(ctx context.Context, qual dax.TableQualifier, sql str orch := newQualifiedOrchestrator(q.orchestrator, qual, q.mds) // Importer - imp := newBatchImporter(idkmds.NewImporter(q.mds, nil), qual, q.mds) + imp := idkmds.NewImporter(q.mds, qual, nil) // TODO(tlt): this obviously doesn't work; we don't have an API here. We // need a dax-compatible implementation of the SystemAPI (or at least a diff --git a/idk/ingest.go b/idk/ingest.go index ac16aa908..a7a272499 100644 --- a/idk/ingest.go +++ b/idk/ingest.go @@ -124,7 +124,7 @@ type Main struct { index *pilosaclient.Index grpcClient *pilosagrpc.GRPCClient - NewImporterFn func() pilosabatch.Importer `flag:"-"` + NewImporterFn func() pilosacore.Importer `flag:"-"` SchemaManager SchemaManager `flag:"-"` Qtbl *dax.QualifiedTable `flag:"-"` @@ -1026,8 +1026,8 @@ func (m *Main) setupClient() (*tls.Config, error) { m.Index = string(qtbl.Key()) m.SchemaManager = mds.NewSchemaManager(dax.Address(m.MDSAddress), qual, m.log) - m.NewImporterFn = func() pilosabatch.Importer { - return mds.NewImporter(mdsClient, qtbl) + m.NewImporterFn = func() pilosacore.Importer { + return mds.NewImporter(mdsClient, qtbl.Qualifier(), &qtbl.Table) } } else { m.SchemaManager = m.client @@ -1859,15 +1859,19 @@ func (m *Main) newBatch(fields []*pilosaclient.Field) (pilosabatch.RecordBatch, pilosabatch.OptKeyTranslateBatchSize(m.KeyTranslateBatchSize), pilosabatch.OptUseShardTransactionalEndpoint(m.UseShardTransactionalEndpoint), } - var importer pilosabatch.Importer + var importer pilosacore.Importer if m.useMDS() { importer = m.NewImporterFn() } else { - importer = pilosaclient.NewImporter(m.client) + sapi := pilosaclient.NewSchemaAPI(m.client) + importer = pilosaclient.NewImporter(m.client, sapi) } opts = append(opts, pilosabatch.OptImporter(importer)) - return pilosabatch.NewBatch(importer, m.BatchSize, pilosaclient.FromClientIndex(m.index), pilosaclient.FromClientFields(fields), opts...) + ii := pilosaclient.FromClientIndex(m.index) + tbl := pilosacore.IndexInfoToTable(ii) + + return pilosabatch.NewBatch(importer, m.BatchSize, tbl, pilosaclient.FromClientFields(fields), opts...) } // validateField ensures that the field is configured correctly. diff --git a/idk/ingest_test.go b/idk/ingest_test.go index c9fe5403a..1776124f3 100644 --- a/idk/ingest_test.go +++ b/idk/ingest_test.go @@ -15,6 +15,7 @@ import ( "time" "github.com/golang-jwt/jwt" + pilosa "github.com/molecula/featurebase/v3" "github.com/molecula/featurebase/v3/authn" batch "github.com/molecula/featurebase/v3/batch" pilosaclient "github.com/molecula/featurebase/v3/client" @@ -56,8 +57,8 @@ func configureTestFlagsMDS(main *Main, address dax.Address, qtbl *dax.QualifiedT main.Index = string(qtbl.Key()) mdsClient := mdsclient.New(dax.Address(address), logger.StderrLogger) - main.NewImporterFn = func() batch.Importer { - return mds.NewImporter(mdsClient, qtbl) + main.NewImporterFn = func() pilosa.Importer { + return mds.NewImporter(mdsClient, qtbl.TableQualifier, &qtbl.Table) } } diff --git a/idk/mds/importer.go b/idk/mds/importer.go index a6a0a6bdd..4f0cc0863 100644 --- a/idk/mds/importer.go +++ b/idk/mds/importer.go @@ -6,7 +6,6 @@ import ( "time" featurebase "github.com/molecula/featurebase/v3" - "github.com/molecula/featurebase/v3/batch" featurebaseclient "github.com/molecula/featurebase/v3/client" "github.com/molecula/featurebase/v3/dax" "github.com/molecula/featurebase/v3/dax/mds/controller/partitioner" @@ -15,20 +14,22 @@ import ( ) // Ensure type implements interface. -var _ batch.Importer = &importer{} +var _ featurebase.Importer = &importer{} // importer type importer struct { mds MDS mu sync.Mutex - qtbl *dax.QualifiedTable + qual dax.TableQualifier + tbl *dax.Table } -func NewImporter(mds MDS, qtbl *dax.QualifiedTable) *importer { +func NewImporter(mds MDS, qual dax.TableQualifier, tbl *dax.Table) *importer { return &importer{ mds: mds, - qtbl: qtbl, + qual: qual, + tbl: tbl, } } @@ -57,8 +58,8 @@ func (m *importer) FinishTransaction(ctx context.Context, id string) (*featureba return nil, nil } -func (m *importer) CreateIndexKeys(ctx context.Context, idx *featurebase.IndexInfo, keys ...string) (map[string]uint64, error) { - qtbl, err := m.getQtbl(ctx, idx.Name) +func (m *importer) CreateTableKeys(ctx context.Context, tid dax.TableID, keys ...string) (map[string]uint64, error) { + qtbl, err := m.getQtbl(ctx, tid) if err != nil { return nil, errors.Wrapf(err, "getting qtbl") } @@ -84,7 +85,8 @@ func (m *importer) CreateIndexKeys(ctx context.Context, idx *featurebase.IndexIn return nil, errors.Wrap(err, "getting featurebase client") } - stringToIDMap, err := fbClient.CreateIndexKeys(featurebaseclient.ToClientIndex(idx), ks...) + cidx := featurebaseclient.QTableToClientIndex(qtbl) + stringToIDMap, err := fbClient.CreateIndexKeys(cidx, ks...) if err != nil { return nil, errors.Wrapf(err, "creating index keys for partition: %d", partition) } @@ -97,8 +99,8 @@ func (m *importer) CreateIndexKeys(ctx context.Context, idx *featurebase.IndexIn return out, nil } -func (m *importer) CreateFieldKeys(ctx context.Context, index string, field *featurebase.FieldInfo, keys ...string) (map[string]uint64, error) { - qtbl, err := m.getQtbl(ctx, index) +func (m *importer) CreateFieldKeys(ctx context.Context, tid dax.TableID, fname dax.FieldName, keys ...string) (map[string]uint64, error) { + qtbl, err := m.getQtbl(ctx, tid) if err != nil { return nil, errors.Wrapf(err, "getting qtbl") } @@ -121,7 +123,7 @@ func (m *importer) CreateFieldKeys(ctx context.Context, index string, field *fea return nil, errors.Wrap(err, "getting featurebase client") } - cfld, err := featurebaseclient.ToClientField(index, field) + cfld, err := featurebaseclient.TableFieldToClientField(qtbl, fname) if err != nil { return nil, errors.Wrap(err, "converting fieldinfo to client field") } @@ -129,8 +131,8 @@ func (m *importer) CreateFieldKeys(ctx context.Context, index string, field *fea return fbClient.CreateFieldKeys(cfld, keys...) } -func (m *importer) ImportRoaringBitmap(ctx context.Context, index string, field *featurebase.FieldInfo, shard uint64, views map[string]*roaring.Bitmap, clear bool) error { - qtbl, err := m.getQtbl(ctx, index) +func (m *importer) ImportRoaringBitmap(ctx context.Context, tid dax.TableID, fld *dax.Field, shard uint64, views map[string]*roaring.Bitmap, clear bool) error { + qtbl, err := m.getQtbl(ctx, tid) if err != nil { return errors.Wrapf(err, "getting qtbl") } @@ -146,7 +148,7 @@ func (m *importer) ImportRoaringBitmap(ctx context.Context, index string, field return errors.Wrap(err, "getting featurebase client") } - cfld, err := featurebaseclient.ToClientField(index, field) + cfld, err := featurebaseclient.TableFieldToClientField(qtbl, fld.Name) if err != nil { return errors.Wrap(err, "converting fieldinfo to client field") } @@ -154,8 +156,8 @@ func (m *importer) ImportRoaringBitmap(ctx context.Context, index string, field return fbClient.ImportRoaringBitmap(cfld, shard, views, clear) } -func (m *importer) ImportRoaringShard(ctx context.Context, index string, shard uint64, request *featurebase.ImportRoaringShardRequest) error { - qtbl, err := m.getQtbl(ctx, index) +func (m *importer) ImportRoaringShard(ctx context.Context, tid dax.TableID, shard uint64, request *featurebase.ImportRoaringShardRequest) error { + qtbl, err := m.getQtbl(ctx, tid) if err != nil { return errors.Wrapf(err, "getting qtbl") } @@ -171,11 +173,11 @@ func (m *importer) ImportRoaringShard(ctx context.Context, index string, shard u return errors.Wrap(err, "getting featurebase client") } - return fbClient.ImportRoaringShard(index, shard, request) + return fbClient.ImportRoaringShard(string(qtbl.Key()), shard, request) } -func (m *importer) EncodeImportValues(ctx context.Context, index string, field *featurebase.FieldInfo, shard uint64, vals []int64, ids []uint64, clear bool) (path string, data []byte, err error) { - qtbl, err := m.getQtbl(ctx, index) +func (m *importer) EncodeImportValues(ctx context.Context, tid dax.TableID, fld *dax.Field, shard uint64, vals []int64, ids []uint64, clear bool) (path string, data []byte, err error) { + qtbl, err := m.getQtbl(ctx, tid) if err != nil { return "", nil, errors.Wrapf(err, "getting qtbl") } @@ -191,7 +193,7 @@ func (m *importer) EncodeImportValues(ctx context.Context, index string, field * return "", nil, errors.Wrap(err, "getting featurebase client") } - cfld, err := featurebaseclient.ToClientField(index, field) + cfld, err := featurebaseclient.TableFieldToClientField(qtbl, fld.Name) if err != nil { return "", nil, errors.Wrap(err, "converting fieldinfo to client field") } @@ -199,8 +201,8 @@ func (m *importer) EncodeImportValues(ctx context.Context, index string, field * return fbClient.EncodeImportValues(cfld, shard, vals, ids, clear) } -func (m *importer) EncodeImport(ctx context.Context, index string, field *featurebase.FieldInfo, shard uint64, vals, ids []uint64, clear bool) (path string, data []byte, err error) { - qtbl, err := m.getQtbl(ctx, index) +func (m *importer) EncodeImport(ctx context.Context, tid dax.TableID, fld *dax.Field, shard uint64, vals, ids []uint64, clear bool) (path string, data []byte, err error) { + qtbl, err := m.getQtbl(ctx, tid) if err != nil { return "", nil, errors.Wrapf(err, "getting qtbl") } @@ -216,7 +218,7 @@ func (m *importer) EncodeImport(ctx context.Context, index string, field *featur return "", nil, errors.Wrap(err, "getting featurebase client") } - cfld, err := featurebaseclient.ToClientField(index, field) + cfld, err := featurebaseclient.TableFieldToClientField(qtbl, fld.Name) if err != nil { return "", nil, errors.Wrap(err, "converting fieldinfo to client field") } @@ -224,8 +226,8 @@ func (m *importer) EncodeImport(ctx context.Context, index string, field *featur return fbClient.EncodeImport(cfld, shard, vals, ids, clear) } -func (m *importer) DoImport(ctx context.Context, index string, field *featurebase.FieldInfo, shard uint64, path string, data []byte) error { - qtbl, err := m.getQtbl(ctx, index) +func (m *importer) DoImport(ctx context.Context, tid dax.TableID, fld *dax.Field, shard uint64, path string, data []byte) error { + qtbl, err := m.getQtbl(ctx, tid) if err != nil { return errors.Wrapf(err, "getting qtbl") } @@ -241,7 +243,7 @@ func (m *importer) DoImport(ctx context.Context, index string, field *featurebas return errors.Wrap(err, "getting featurebase client") } - return fbClient.DoImport(index, shard, path, data) + return fbClient.DoImport(string(qtbl.Key()), shard, path, data) } func (m *importer) StatsTiming(name string, value time.Duration, rate float64) {} @@ -254,23 +256,22 @@ func (m *importer) StatsTiming(name string, value time.Duration, rate float64) { // yet. So this method allows us to use the table which is passed into each // method to determine the table. We look it up from mds schema once and save it // in m.qtbl for any further method calls. -func (m *importer) getQtbl(ctx context.Context, table string) (*dax.QualifiedTable, error) { +func (m *importer) getQtbl(ctx context.Context, tid dax.TableID) (*dax.QualifiedTable, error) { m.mu.Lock() defer m.mu.Unlock() - if m.qtbl != nil { - return m.qtbl, nil + if m.tbl != nil { + return dax.NewQualifiedTable(m.qual, m.tbl), nil } - tkey := dax.TableKey(table) - qtid := tkey.QualifiedTableID() + qtid := dax.NewQualifiedTableID(m.qual, tid) qtbl, err := m.mds.Table(ctx, qtid) if err != nil { return nil, errors.Wrap(err, "getting table") } - m.qtbl = qtbl + m.tbl = &qtbl.Table return qtbl, nil } diff --git a/importer.go b/importer.go new file mode 100644 index 000000000..9bfa3b79d --- /dev/null +++ b/importer.go @@ -0,0 +1,88 @@ +package pilosa + +import ( + "context" + "time" + + "github.com/molecula/featurebase/v3/dax" + "github.com/molecula/featurebase/v3/roaring" +) + +type Importer interface { + StartTransaction(ctx context.Context, id string, timeout time.Duration, exclusive bool, requestTimeout time.Duration) (*Transaction, error) + FinishTransaction(ctx context.Context, id string) (*Transaction, error) + CreateTableKeys(ctx context.Context, tid dax.TableID, keys ...string) (map[string]uint64, error) + CreateFieldKeys(ctx context.Context, tid dax.TableID, fname dax.FieldName, keys ...string) (map[string]uint64, error) + ImportRoaringBitmap(ctx context.Context, tid dax.TableID, fld *dax.Field, shard uint64, views map[string]*roaring.Bitmap, clear bool) error + ImportRoaringShard(ctx context.Context, tid dax.TableID, shard uint64, request *ImportRoaringShardRequest) error + EncodeImportValues(ctx context.Context, tid dax.TableID, fld *dax.Field, shard uint64, vals []int64, ids []uint64, clear bool) (path string, data []byte, err error) + EncodeImport(ctx context.Context, tid dax.TableID, fld *dax.Field, shard uint64, vals, ids []uint64, clear bool) (path string, data []byte, err error) + DoImport(ctx context.Context, tid dax.TableID, fld *dax.Field, shard uint64, path string, data []byte) error + + StatsTiming(name string, value time.Duration, rate float64) +} + +// Ensure type implements interface. +var _ Importer = &onPremImporter{} + +// onPremImporter is a wrapper around API which implements the Importer +// interface. This is currently only used by sql3 running locally in standard +// (i.e not "serverless") mode. Because sql3 always sets +// `useShardTransactionalEndpoint = true`, There are several methods which this +// implemtation of the Importer interface does not use, and therefore they +// intentionally no-op. +type onPremImporter struct { + api *API +} + +func NewOnPremImporter(api *API) *onPremImporter { + return &onPremImporter{ + api: api, + } +} + +func (i *onPremImporter) StartTransaction(ctx context.Context, id string, timeout time.Duration, exclusive bool, requestTimeout time.Duration) (*Transaction, error) { + return i.api.StartTransaction(ctx, id, timeout, exclusive, false) +} + +func (i *onPremImporter) FinishTransaction(ctx context.Context, id string) (*Transaction, error) { + return i.api.FinishTransaction(ctx, id, false) +} + +func (i *onPremImporter) CreateTableKeys(ctx context.Context, tid dax.TableID, keys ...string) (map[string]uint64, error) { + return i.api.CreateIndexKeys(ctx, string(tid), keys...) +} + +func (i *onPremImporter) CreateFieldKeys(ctx context.Context, tid dax.TableID, fname dax.FieldName, keys ...string) (map[string]uint64, error) { + return i.api.CreateFieldKeys(ctx, string(tid), string(fname), keys...) +} + +func (i *onPremImporter) ImportRoaringBitmap(ctx context.Context, tid dax.TableID, fld *dax.Field, shard uint64, views map[string]*roaring.Bitmap, clear bool) error { + // This intentionally no-ops. See comment on struct. + return nil +} + +func (i *onPremImporter) ImportRoaringShard(ctx context.Context, tid dax.TableID, shard uint64, request *ImportRoaringShardRequest) error { + return i.api.ImportRoaringShard(ctx, string(tid), shard, request) +} + +// EncodeImportValues is kind of weird. We're trying to mimic what the client +// does here (because the Importer interface was originally based off of the +// client methods). So we end up generating a protobuf-encode byte slice. And we +// don't really use path. +func (i *onPremImporter) EncodeImportValues(ctx context.Context, tid dax.TableID, fld *dax.Field, shard uint64, vals []int64, ids []uint64, clear bool) (path string, data []byte, err error) { + // This intentionally no-ops. See comment on struct. + return "", nil, nil +} + +func (i *onPremImporter) EncodeImport(ctx context.Context, tid dax.TableID, fld *dax.Field, shard uint64, vals, ids []uint64, clear bool) (path string, data []byte, err error) { + // This intentionally no-ops. See comment on struct. + return "", nil, nil +} + +func (i *onPremImporter) DoImport(ctx context.Context, _ dax.TableID, fld *dax.Field, shard uint64, path string, data []byte) error { + // This intentionally no-ops. See comment on struct. + return nil +} + +func (i *onPremImporter) StatsTiming(name string, value time.Duration, rate float64) {} diff --git a/schema.go b/schema.go index 399ab9a2b..562e8226d 100644 --- a/schema.go +++ b/schema.go @@ -141,17 +141,11 @@ func IndexInfoToTable(ii *IndexInfo) *dax.Table { CreatedAt: ii.CreatedAt, } - // // sortedFields will contain the sorted list of fields from IndexInfo. - // sortedFields := make([]*FieldInfo, 0, len(ii.Fields)) - // Sort ii.Fields by CreatedAt before adding them to sortedFields. sort.Slice(ii.Fields, func(i, j int) bool { return ii.Fields[i].CreatedAt < ii.Fields[j].CreatedAt }) - // // Add the sorted fields to sortedFields. - // sortedFields = append(sortedFields, ii.Fields...) - // Add the _id Field. var idType dax.BaseType = dax.BaseTypeID if ii.Options.Keys { diff --git a/server/server.go b/server/server.go index 6864b66ac..f2d66d5cb 100644 --- a/server/server.go +++ b/server/server.go @@ -29,7 +29,6 @@ import ( pilosa "github.com/molecula/featurebase/v3" "github.com/molecula/featurebase/v3/authn" "github.com/molecula/featurebase/v3/authz" - "github.com/molecula/featurebase/v3/batch" "github.com/molecula/featurebase/v3/boltdb" "github.com/molecula/featurebase/v3/dax" "github.com/molecula/featurebase/v3/dax/computer" @@ -569,9 +568,9 @@ func (m *Command) setupServer() error { executionPlannerFn := func(e pilosa.Executor, api *pilosa.API, sql string) sql3.CompilePlanner { fapi := pilosa.NewOnPremSchema(api) fsapi := &pilosa.FeatureBaseSystemAPI{API: api} - fimp := &batch.FeaturebaseImporter{API: api} + imp := pilosa.NewOnPremImporter(api) - return planner.NewExecutionPlanner(e, fapi, fsapi, m.Server.SystemLayer, fimp, m.logger, sql) + return planner.NewExecutionPlanner(e, fapi, fsapi, m.Server.SystemLayer, imp, m.logger, sql) } serverOptions := []pilosa.ServerOption{ diff --git a/sql3/planner/executionplanner.go b/sql3/planner/executionplanner.go index 605269ed4..d8c7e40cf 100644 --- a/sql3/planner/executionplanner.go +++ b/sql3/planner/executionplanner.go @@ -6,7 +6,6 @@ import ( "context" pilosa "github.com/molecula/featurebase/v3" - "github.com/molecula/featurebase/v3/batch" "github.com/molecula/featurebase/v3/logger" "github.com/molecula/featurebase/v3/sql3" "github.com/molecula/featurebase/v3/sql3/parser" @@ -27,13 +26,13 @@ type ExecutionPlanner struct { schemaAPI pilosa.SchemaAPI systemAPI pilosa.SystemAPI systemLayerAPI pilosa.SystemLayerAPI - importer batch.Importer + importer pilosa.Importer logger logger.Logger sql string scopeStack *scopeStack } -func NewExecutionPlanner(executor pilosa.Executor, schemaAPI pilosa.SchemaAPI, systemAPI pilosa.SystemAPI, systemLayerAPI pilosa.SystemLayerAPI, importer batch.Importer, logger logger.Logger, sql string) *ExecutionPlanner { +func NewExecutionPlanner(executor pilosa.Executor, schemaAPI pilosa.SchemaAPI, systemAPI pilosa.SystemAPI, systemLayerAPI pilosa.SystemLayerAPI, importer pilosa.Importer, logger logger.Logger, sql string) *ExecutionPlanner { return &ExecutionPlanner{ executor: executor, schemaAPI: newSystemTableDefintionsWrapper(schemaAPI), diff --git a/sql3/planner/opinsert.go b/sql3/planner/opinsert.go index 42bf0c748..31503ec42 100644 --- a/sql3/planner/opinsert.go +++ b/sql3/planner/opinsert.go @@ -173,7 +173,7 @@ func (i *insertRowIter) Next(ctx context.Context) (types.Row, error) { counter++ } - batch, err := fbbatch.NewBatch(i.planner.importer, batchSize, idxInfo, idxInfo.Fields, + batch, err := fbbatch.NewBatch(i.planner.importer, batchSize, tbl, idxInfo.Fields, fbbatch.OptUseShardTransactionalEndpoint(true), ) if err != nil {