mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-08-28 02:44:59 +00:00
* First pass at RetryWithTx * Refactor RetryWithTx to take a writable bool (instead of reads, writes) * Implment DirectiveMethodDiff This commit adds support for a Directive to contain only the diffs (as opposed to the full Directive). * Update controller tests to allow for DirectiveMethodDiff (over Full) * Update RetryWithTx to retry on duplicate key constraint. If two concurrent processes call IngestShard() for the same shard, both were trying to insert the same job into the jobs table. That resulted in a duplicate key error from the database. We want to include that error in the list of errors for which RetryWithTx should retry. * Remove unused method: Directive.TranslatePartitions() * Replace query in a loop with a single query We had a query which was looking to see if a job already existed. That query was inside a loop, and could potentially generate 256 queries (for example). This commit replaces that logic so that we use a single query wiht an `IN ()` clause. * Convert to directive version-by-address This commit uses a separate directive version per address. It moves the version get/increment back inside the buildDirective method so that if two concurrent processes are building a directive for the same address, one of them will get rolled back trying to commit the version update. * Migration for directive version by address * Add a comment about DirectiveVersion lock/unlock logic * Remove AddLastWins * fix linter * handle error in walkdir * fix test failures from removing AddLastWins
280 lines
9 KiB
Go
280 lines
9 KiB
Go
package serverless
|
|
|
|
import (
|
|
"context"
|
|
"sync"
|
|
"time"
|
|
|
|
featurebase "github.com/featurebasedb/featurebase/v3"
|
|
fbclient "github.com/featurebasedb/featurebase/v3/client"
|
|
"github.com/featurebasedb/featurebase/v3/dax"
|
|
"github.com/featurebasedb/featurebase/v3/dax/controller/partitioner"
|
|
"github.com/featurebasedb/featurebase/v3/roaring"
|
|
"github.com/pkg/errors"
|
|
)
|
|
|
|
// Ensure type implements interface.
|
|
var _ featurebase.Importer = &importer{}
|
|
|
|
// importer
|
|
type importer struct {
|
|
controller dax.Controller
|
|
|
|
mu sync.Mutex
|
|
qdbid dax.QualifiedDatabaseID
|
|
tbl *dax.Table
|
|
}
|
|
|
|
func NewImporter(controller dax.Controller, qdbid dax.QualifiedDatabaseID, tbl *dax.Table) *importer {
|
|
return &importer{
|
|
controller: controller,
|
|
qdbid: qdbid,
|
|
tbl: tbl,
|
|
}
|
|
}
|
|
|
|
var fbClientCache = map[dax.Address]*fbclient.Client{}
|
|
var fbClientCacheMu sync.Mutex
|
|
|
|
func (m *importer) fbClient(address dax.Address) (*fbclient.Client, error) {
|
|
fbClientCacheMu.Lock()
|
|
defer fbClientCacheMu.Unlock()
|
|
client := fbClientCache[address]
|
|
if client != nil {
|
|
return client, nil
|
|
}
|
|
client, err := fbclient.NewClient(address.HostPort(),
|
|
fbclient.OptClientRetries(2),
|
|
fbclient.OptClientTotalPoolSize(1000),
|
|
fbclient.OptClientPoolSizePerRoute(400),
|
|
fbclient.OptClientPathPrefix(address.Path()),
|
|
//fbclient.OptClientStatsClient(m.stats),
|
|
)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
fbClientCache[address] = client
|
|
return client, nil
|
|
}
|
|
|
|
func (m *importer) StartTransaction(ctx context.Context, id string, timeout time.Duration, exclusive bool, requestTimeout time.Duration) (*featurebase.Transaction, error) {
|
|
return &featurebase.Transaction{
|
|
ID: "not-used",
|
|
}, nil
|
|
}
|
|
|
|
func (m *importer) FinishTransaction(ctx context.Context, id string) (*featurebase.Transaction, error) {
|
|
return nil, nil
|
|
}
|
|
|
|
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")
|
|
}
|
|
|
|
out := make(map[string]uint64)
|
|
|
|
partitioner := partitioner.NewPartitioner()
|
|
|
|
// Get the partition for each key (map[int][]string).
|
|
partitions := partitioner.PartitionsForKeys(qtbl.Key(), qtbl.PartitionN, keys...)
|
|
|
|
// TODO: we can be more efficient here by calling IngestPartitions() with
|
|
// all the partitions at once, then getting the distinct list of addresses
|
|
// and looping over that instead.
|
|
for partition, ks := range partitions {
|
|
address, err := m.controller.IngestPartition(context.Background(), qtbl.QualifiedID(), partition)
|
|
if err != nil {
|
|
return nil, errors.Wrapf(err, "calling ingest-partition on table: %s, partition: %d", qtbl, partition)
|
|
}
|
|
|
|
fbClient, err := m.fbClient(address)
|
|
if err != nil {
|
|
return nil, errors.Wrap(err, "getting featurebase client")
|
|
}
|
|
|
|
cidx := fbclient.QTableToClientIndex(qtbl)
|
|
stringToIDMap, err := fbClient.CreateIndexKeys(cidx, ks...)
|
|
if err != nil {
|
|
return nil, errors.Wrapf(err, "creating index keys for partition: %d", partition)
|
|
}
|
|
|
|
for str, id := range stringToIDMap {
|
|
out[str] = id
|
|
}
|
|
}
|
|
|
|
return out, nil
|
|
}
|
|
|
|
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")
|
|
}
|
|
|
|
// For now, we are going to direct all field key translation to the same
|
|
// node handling index key translation for the table, partition 0.
|
|
// TODO: we should be able to partition field key translation on fieldName.
|
|
// If we do that, we might also consider whether we want to support a
|
|
// different partitionN for field translation.
|
|
partition := dax.PartitionNum(0)
|
|
|
|
address, err := m.controller.IngestPartition(context.Background(), qtbl.QualifiedID(), partition)
|
|
if err != nil {
|
|
return nil, errors.Wrapf(err, "calling ingest-partition on table: %s, partition: %d", qtbl, partition)
|
|
}
|
|
|
|
// Set up a FeatureBase client with address.
|
|
fbClient, err := m.fbClient(address)
|
|
if err != nil {
|
|
return nil, errors.Wrap(err, "getting featurebase client")
|
|
}
|
|
|
|
cfld, err := fbclient.TableFieldToClientField(qtbl, fname)
|
|
if err != nil {
|
|
return nil, errors.Wrap(err, "converting fieldinfo to client field")
|
|
}
|
|
|
|
return fbClient.CreateFieldKeys(cfld, keys...)
|
|
}
|
|
|
|
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")
|
|
}
|
|
|
|
address, err := m.controller.IngestShard(context.Background(), qtbl.QualifiedID(), dax.ShardNum(shard))
|
|
if err != nil {
|
|
return errors.Wrap(err, "calling ingest-shard")
|
|
}
|
|
|
|
// Set up a FeatureBase client with address.
|
|
fbClient, err := m.fbClient(address)
|
|
if err != nil {
|
|
return errors.Wrap(err, "getting featurebase client")
|
|
}
|
|
|
|
cfld, err := fbclient.TableFieldToClientField(qtbl, fld.Name)
|
|
if err != nil {
|
|
return errors.Wrap(err, "converting fieldinfo to client field")
|
|
}
|
|
|
|
return fbClient.ImportRoaringBitmap(cfld, shard, views, clear)
|
|
}
|
|
|
|
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")
|
|
}
|
|
|
|
address, err := m.controller.IngestShard(context.Background(), qtbl.QualifiedID(), dax.ShardNum(shard))
|
|
if err != nil {
|
|
return errors.Wrap(err, "calling ingest-shard")
|
|
}
|
|
|
|
// Set up a FeatureBase client with address.
|
|
fbClient, err := m.fbClient(address)
|
|
if err != nil {
|
|
return errors.Wrap(err, "getting featurebase client")
|
|
}
|
|
|
|
return fbClient.ImportRoaringShard(string(qtbl.Key()), shard, request)
|
|
}
|
|
|
|
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")
|
|
}
|
|
|
|
// Since we're calling EncodeImportValues on the client, we don't actually
|
|
// need a valid client (that method doesn't actually use the client).
|
|
// Really, that method should be a function on the client package rather
|
|
// than a method on the Client type.
|
|
fbClient, err := m.fbClient("")
|
|
if err != nil {
|
|
return "", nil, errors.Wrap(err, "getting featurebase client")
|
|
}
|
|
|
|
cfld, err := fbclient.TableFieldToClientField(qtbl, fld.Name)
|
|
if err != nil {
|
|
return "", nil, errors.Wrap(err, "converting fieldinfo to client field")
|
|
}
|
|
|
|
return fbClient.EncodeImportValues(cfld, shard, vals, ids, clear)
|
|
}
|
|
|
|
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")
|
|
}
|
|
|
|
// Since we're calling EncodeImportValues on the client, we don't actually
|
|
// need a valid client (that method doesn't actually use the client).
|
|
// Really, that method should be a function on the client package rather
|
|
// than a method on the Client type.
|
|
fbClient, err := m.fbClient("")
|
|
if err != nil {
|
|
return "", nil, errors.Wrap(err, "getting featurebase client")
|
|
}
|
|
|
|
cfld, err := fbclient.TableFieldToClientField(qtbl, fld.Name)
|
|
if err != nil {
|
|
return "", nil, errors.Wrap(err, "converting fieldinfo to client field")
|
|
}
|
|
|
|
return fbClient.EncodeImport(cfld, shard, vals, ids, clear)
|
|
}
|
|
|
|
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")
|
|
}
|
|
|
|
address, err := m.controller.IngestShard(context.Background(), qtbl.QualifiedID(), dax.ShardNum(shard))
|
|
if err != nil {
|
|
return errors.Wrap(err, "calling ingest-shard")
|
|
}
|
|
|
|
// Set up a FeatureBase client with address.
|
|
fbClient, err := m.fbClient(address)
|
|
if err != nil {
|
|
return errors.Wrap(err, "getting featurebase client")
|
|
}
|
|
|
|
return fbClient.DoImport(string(qtbl.Key()), shard, path, data)
|
|
}
|
|
|
|
// getQtbl takes a table (TableKey) and sets the local m.qtbl value. When we
|
|
// originally set up this type, it was only used by IDK, and the table was known
|
|
// at the beginning of the process, so it could be set on this import. But
|
|
// later, we used this importer in the Queryer, and it gets set up prior to
|
|
// parsing the table out of sql; which means we don't know what the table is
|
|
// 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, tid dax.TableID) (*dax.QualifiedTable, error) {
|
|
m.mu.Lock()
|
|
defer m.mu.Unlock()
|
|
|
|
if m.tbl != nil {
|
|
return dax.NewQualifiedTable(m.qdbid, m.tbl), nil
|
|
}
|
|
|
|
qtid := dax.NewQualifiedTableID(m.qdbid, tid)
|
|
|
|
qtbl, err := m.controller.TableByID(ctx, qtid)
|
|
if err != nil {
|
|
return nil, errors.Wrap(err, "getting table")
|
|
}
|
|
|
|
m.tbl = &qtbl.Table
|
|
|
|
return qtbl, nil
|
|
}
|