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.
This commit is contained in:
Travis Turner 2022-12-09 14:24:47 -06:00 committed by GitHub
parent 57ce7c4c0e
commit 12d608c80d
No known key found for this signature in database
GPG key ID: 4AEE18F83AFDEB23
15 changed files with 401 additions and 504 deletions

View file

@ -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)
})

View file

@ -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)
}

View file

@ -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) {}

View file

@ -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) {

View file

@ -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)
}

View file

@ -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
}

View file

@ -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

View file

@ -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.

View file

@ -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)
}
}

View file

@ -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
}

88
importer.go Normal file
View file

@ -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) {}

View file

@ -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 {

View file

@ -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{

View file

@ -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),

View file

@ -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 {