mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-09-12 07:41:02 +00:00
Make interfaces more specific than "MDS" (#2352)
* Make interfaces more specific than "MDS" - Introduce `dax.Schemar` interface - Introduce `dax.Noder` interface - The rest is generally to standardize on the new interfaces. - Remove `pilosa.SchemaInfoAPI` interface - Move `TranslateNode` and `ComputeNode` types from controller to dax package - Remove `queryer.FeatureBaseImporter` - Remove `queryer.MDS` interface - Remove `queryer.Importer` interface - Identify types using an "MDS" interface and split into Noder/Schemar as necessary - Changed `Queryer.orchestrator` to a `map[qual]*qualifiedOrchestrator` because we can't share an orchestrator across quals * Convert orchestrator to use TableKeyer
This commit is contained in:
parent
63cfdb5078
commit
14f1930004
20 changed files with 483 additions and 511 deletions
5
api.go
5
api.go
|
|
@ -3385,11 +3385,6 @@ type SchemaAPI interface {
|
|||
DeleteField(ctx context.Context, tname dax.TableName, fname dax.FieldName) error
|
||||
}
|
||||
|
||||
type SchemaInfoAPI interface {
|
||||
IndexInfo(ctx context.Context, indexName string) (*IndexInfo, error)
|
||||
FieldInfo(ctx context.Context, indexName, fieldName string) (*FieldInfo, error)
|
||||
}
|
||||
|
||||
type ClusterNode struct {
|
||||
ID string
|
||||
State string
|
||||
|
|
|
|||
|
|
@ -10,7 +10,6 @@ import (
|
|||
"net/http"
|
||||
|
||||
"github.com/molecula/featurebase/v3/dax"
|
||||
"github.com/molecula/featurebase/v3/dax/mds/controller"
|
||||
mdshttp "github.com/molecula/featurebase/v3/dax/mds/http"
|
||||
"github.com/molecula/featurebase/v3/errors"
|
||||
"github.com/molecula/featurebase/v3/logger"
|
||||
|
|
@ -49,6 +48,20 @@ func (c *Client) Health() bool {
|
|||
return true
|
||||
}
|
||||
|
||||
// TODO(tlt): collapse Table into this
|
||||
func (c *Client) TableByID(ctx context.Context, qtid dax.QualifiedTableID) (*dax.QualifiedTable, error) {
|
||||
return c.Table(ctx, qtid)
|
||||
}
|
||||
|
||||
// TODO(tlt): collapse TableID into this
|
||||
func (c *Client) TableByName(ctx context.Context, qual dax.TableQualifier, tname dax.TableName) (*dax.QualifiedTable, error) {
|
||||
qtid, err := c.TableID(ctx, qual, tname)
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "getting table id")
|
||||
}
|
||||
return c.Table(ctx, qtid)
|
||||
}
|
||||
|
||||
func (c *Client) Table(ctx context.Context, qtid dax.QualifiedTableID) (*dax.QualifiedTable, error) {
|
||||
url := fmt.Sprintf("%s/table", c.address.WithScheme(defaultScheme))
|
||||
|
||||
|
|
@ -336,11 +349,11 @@ func (c *Client) IngestPartition(ctx context.Context, qtid dax.QualifiedTableID,
|
|||
return isr.Address, nil
|
||||
}
|
||||
|
||||
func (c *Client) ComputeNodes(ctx context.Context, qtid dax.QualifiedTableID, shards ...dax.ShardNum) ([]controller.ComputeNode, error) {
|
||||
func (c *Client) ComputeNodes(ctx context.Context, qtid dax.QualifiedTableID, shards ...dax.ShardNum) ([]dax.ComputeNode, error) {
|
||||
url := fmt.Sprintf("%s/compute-nodes", c.address.WithScheme(defaultScheme))
|
||||
c.logger.Debugf("ComputeNodes url: %s", url)
|
||||
|
||||
var nodes []controller.ComputeNode
|
||||
var nodes []dax.ComputeNode
|
||||
|
||||
req := &mdshttp.ComputeNodesRequest{
|
||||
Table: qtid,
|
||||
|
|
@ -374,11 +387,11 @@ func (c *Client) ComputeNodes(ctx context.Context, qtid dax.QualifiedTableID, sh
|
|||
return cnr.ComputeNodes, nil
|
||||
}
|
||||
|
||||
func (c *Client) TranslateNodes(ctx context.Context, qtid dax.QualifiedTableID, partitions ...dax.PartitionNum) ([]controller.TranslateNode, error) {
|
||||
func (c *Client) TranslateNodes(ctx context.Context, qtid dax.QualifiedTableID, partitions ...dax.PartitionNum) ([]dax.TranslateNode, error) {
|
||||
url := fmt.Sprintf("%s/translate-nodes", c.address.WithScheme(defaultScheme))
|
||||
c.logger.Debugf("TranslateNodes url: %s", url)
|
||||
|
||||
var nodes []controller.TranslateNode
|
||||
var nodes []dax.TranslateNode
|
||||
|
||||
req := &mdshttp.TranslateNodesRequest{
|
||||
Table: qtid,
|
||||
|
|
|
|||
|
|
@ -1647,7 +1647,7 @@ func (c *Controller) SnapshotFieldKeys(ctx context.Context, qtid dax.QualifiedTa
|
|||
|
||||
/////////////
|
||||
|
||||
func (c *Controller) ComputeNodes(ctx context.Context, qtid dax.QualifiedTableID, shards dax.ShardNums, isWrite bool) ([]ComputeNode, error) {
|
||||
func (c *Controller) ComputeNodes(ctx context.Context, qtid dax.QualifiedTableID, shards dax.ShardNums, isWrite bool) ([]dax.ComputeNode, error) {
|
||||
inRole := &dax.ComputeRole{
|
||||
TableKey: qtid.Key(),
|
||||
Shards: dax.NewVersionedShards(shards...),
|
||||
|
|
@ -1658,7 +1658,7 @@ func (c *Controller) ComputeNodes(ctx context.Context, qtid dax.QualifiedTableID
|
|||
return nil, errors.Wrap(err, "getting compute nodes")
|
||||
}
|
||||
|
||||
computeNodes := make([]ComputeNode, 0)
|
||||
computeNodes := make([]dax.ComputeNode, 0)
|
||||
|
||||
for _, node := range nodes {
|
||||
role, ok := node.Role.(*dax.ComputeRole)
|
||||
|
|
@ -1668,7 +1668,7 @@ func (c *Controller) ComputeNodes(ctx context.Context, qtid dax.QualifiedTableID
|
|||
return nil, NewErrInternal("not a compute node")
|
||||
}
|
||||
|
||||
computeNodes = append(computeNodes, ComputeNode{
|
||||
computeNodes = append(computeNodes, dax.ComputeNode{
|
||||
Address: node.Address,
|
||||
Table: role.TableKey,
|
||||
Shards: role.Shards.Nums(),
|
||||
|
|
@ -1678,7 +1678,7 @@ func (c *Controller) ComputeNodes(ctx context.Context, qtid dax.QualifiedTableID
|
|||
return computeNodes, nil
|
||||
}
|
||||
|
||||
func (c *Controller) TranslateNodes(ctx context.Context, qtid dax.QualifiedTableID, partitions dax.PartitionNums, isWrite bool) ([]TranslateNode, error) {
|
||||
func (c *Controller) TranslateNodes(ctx context.Context, qtid dax.QualifiedTableID, partitions dax.PartitionNums, isWrite bool) ([]dax.TranslateNode, error) {
|
||||
inRole := &dax.TranslateRole{
|
||||
TableKey: qtid.Key(),
|
||||
Partitions: dax.NewVersionedPartitions(partitions...),
|
||||
|
|
@ -1689,7 +1689,7 @@ func (c *Controller) TranslateNodes(ctx context.Context, qtid dax.QualifiedTable
|
|||
return nil, errors.Wrap(err, "getting translate nodes")
|
||||
}
|
||||
|
||||
translateNodes := make([]TranslateNode, 0)
|
||||
translateNodes := make([]dax.TranslateNode, 0)
|
||||
|
||||
for _, node := range nodes {
|
||||
role, ok := node.Role.(*dax.TranslateRole)
|
||||
|
|
@ -1699,7 +1699,7 @@ func (c *Controller) TranslateNodes(ctx context.Context, qtid dax.QualifiedTable
|
|||
return nil, NewErrInternal("not a translate node")
|
||||
}
|
||||
|
||||
translateNodes = append(translateNodes, TranslateNode{
|
||||
translateNodes = append(translateNodes, dax.TranslateNode{
|
||||
Address: node.Address,
|
||||
Table: role.TableKey,
|
||||
Partitions: role.Partitions.Nums(),
|
||||
|
|
|
|||
|
|
@ -1,19 +0,0 @@
|
|||
package controller
|
||||
|
||||
import "github.com/molecula/featurebase/v3/dax"
|
||||
|
||||
// ComputeNode represents a compute node and the table/shards for which it is
|
||||
// responsible.
|
||||
type ComputeNode struct {
|
||||
Address dax.Address `json:"address"`
|
||||
Table dax.TableKey `json:"table"`
|
||||
Shards dax.ShardNums `json:"shards"`
|
||||
}
|
||||
|
||||
// TranslateNode represents a translate node and the table/partitions for which
|
||||
// it is responsible.
|
||||
type TranslateNode struct {
|
||||
Address dax.Address `json:"address"`
|
||||
Table dax.TableKey `json:"table"`
|
||||
Partitions dax.PartitionNums `json:"partitions"`
|
||||
}
|
||||
|
|
@ -7,7 +7,6 @@ import (
|
|||
"github.com/gorilla/mux"
|
||||
"github.com/molecula/featurebase/v3/dax"
|
||||
"github.com/molecula/featurebase/v3/dax/mds"
|
||||
"github.com/molecula/featurebase/v3/dax/mds/controller"
|
||||
)
|
||||
|
||||
func Handler(mds *mds.MDS) http.Handler {
|
||||
|
|
@ -613,7 +612,7 @@ type ComputeNodesRequest struct {
|
|||
// provided are not included in this response. That might happen if there are
|
||||
// currently no active compute nodes.
|
||||
type ComputeNodesResponse struct {
|
||||
ComputeNodes []controller.ComputeNode `json:"compute-nodes"`
|
||||
ComputeNodes []dax.ComputeNode `json:"compute-nodes"`
|
||||
}
|
||||
|
||||
// POST /translate-nodes
|
||||
|
|
@ -660,5 +659,5 @@ type TranslateNodesRequest struct {
|
|||
// that partitions provided are not included in this response. That might happen
|
||||
// if there are currently no active translate nodes.
|
||||
type TranslateNodesResponse struct {
|
||||
TranslateNodes []controller.TranslateNode `json:"translate-nodes"`
|
||||
TranslateNodes []dax.TranslateNode `json:"translate-nodes"`
|
||||
}
|
||||
|
|
|
|||
|
|
@ -462,7 +462,7 @@ func (m *MDS) DeregisterNodes(ctx context.Context, addrs ...dax.Address) error {
|
|||
|
||||
// ComputeNodes gets the compute nodes responsible for the table/shards
|
||||
// specified in the ComputeNodeRequest.
|
||||
func (m *MDS) ComputeNodes(ctx context.Context, qtid dax.QualifiedTableID, shardNums ...dax.ShardNum) ([]controller.ComputeNode, error) {
|
||||
func (m *MDS) ComputeNodes(ctx context.Context, qtid dax.QualifiedTableID, shardNums ...dax.ShardNum) ([]dax.ComputeNode, error) {
|
||||
if err := m.sanitizeQTID(ctx, &qtid); err != nil {
|
||||
return nil, errors.Wrap(err, "sanitizing")
|
||||
}
|
||||
|
|
@ -476,7 +476,7 @@ func (m *MDS) DebugNodes(ctx context.Context) ([]*dax.Node, error) {
|
|||
|
||||
// TranslateNodes gets the translate nodes responsible for the table/partitions
|
||||
// specified in the TranslateNodeRequest.
|
||||
func (m *MDS) TranslateNodes(ctx context.Context, qtid dax.QualifiedTableID, partitionNums ...dax.PartitionNum) ([]controller.TranslateNode, error) {
|
||||
func (m *MDS) TranslateNodes(ctx context.Context, qtid dax.QualifiedTableID, partitionNums ...dax.PartitionNum) ([]dax.TranslateNode, error) {
|
||||
if err := m.sanitizeQTID(ctx, &qtid); err != nil {
|
||||
return nil, errors.Wrap(err, "sanitizing")
|
||||
}
|
||||
|
|
|
|||
56
dax/node.go
56
dax/node.go
|
|
@ -29,6 +29,62 @@ type NodeService interface {
|
|||
Nodes(context.Context) ([]*Node, error)
|
||||
}
|
||||
|
||||
// ComputeNode represents a compute node and the table/shards for which it is
|
||||
// responsible.
|
||||
type ComputeNode struct {
|
||||
Address Address `json:"address"`
|
||||
Table TableKey `json:"table"`
|
||||
Shards ShardNums `json:"shards"`
|
||||
}
|
||||
|
||||
// TranslateNode represents a translate node and the table/partitions for which
|
||||
// it is responsible.
|
||||
type TranslateNode struct {
|
||||
Address Address `json:"address"`
|
||||
Table TableKey `json:"table"`
|
||||
Partitions PartitionNums `json:"partitions"`
|
||||
}
|
||||
|
||||
type Noder interface {
|
||||
ComputeNodes(ctx context.Context, qtid QualifiedTableID, shards ...ShardNum) ([]ComputeNode, error)
|
||||
TranslateNodes(ctx context.Context, qtid QualifiedTableID, partitions ...PartitionNum) ([]TranslateNode, error)
|
||||
|
||||
// IngestPartition is effectively the "write" version of TranslateNodes. Its
|
||||
// implementations will return the same Address that TranslateNodes would,
|
||||
// but it includes the logic to create/assign the partition if it is not
|
||||
// already being handled by a computer.
|
||||
IngestPartition(ctx context.Context, qtid QualifiedTableID, partition PartitionNum) (Address, error)
|
||||
|
||||
// IngestShard is effectively the "write" version of ComputeNodes. Its
|
||||
// implementations will return the same Address that ComputeNodes would, but
|
||||
// it includes the logic to create/assign the shard if it is not already
|
||||
// being handled by a computer.
|
||||
IngestShard(ctx context.Context, qtid QualifiedTableID, shard ShardNum) (Address, error)
|
||||
}
|
||||
|
||||
// Ensure type implements interface.
|
||||
var _ Noder = &nopNoder{}
|
||||
|
||||
// NopMDS is a no-op implementation of the MDS interface.
|
||||
type nopNoder struct{}
|
||||
|
||||
func NewNopNoder() *nopNoder {
|
||||
return &nopNoder{}
|
||||
}
|
||||
|
||||
func (n *nopNoder) ComputeNodes(ctx context.Context, qtid QualifiedTableID, shards ...ShardNum) ([]ComputeNode, error) {
|
||||
return nil, nil
|
||||
}
|
||||
func (n *nopNoder) IngestPartition(ctx context.Context, qtid QualifiedTableID, partition PartitionNum) (Address, error) {
|
||||
return "", nil
|
||||
}
|
||||
func (n *nopNoder) IngestShard(ctx context.Context, qtid QualifiedTableID, shard ShardNum) (Address, error) {
|
||||
return "", nil
|
||||
}
|
||||
func (n *nopNoder) TranslateNodes(ctx context.Context, qtid QualifiedTableID, partitions ...PartitionNum) ([]TranslateNode, error) {
|
||||
return nil, nil
|
||||
}
|
||||
|
||||
////////////////////////////////////////////////////
|
||||
// Errors
|
||||
////////////////////////////////////////////////////
|
||||
|
|
|
|||
|
|
@ -1,41 +0,0 @@
|
|||
package queryer
|
||||
|
||||
import (
|
||||
"context"
|
||||
|
||||
featurebase "github.com/molecula/featurebase/v3"
|
||||
)
|
||||
|
||||
// Ensure type implements interface.
|
||||
var _ Importer = &FeatureBaseImporter{}
|
||||
|
||||
// FeatureBaseImporter is an implementation of the Importer interface which uses
|
||||
// a pointer to a featurebase.API to make the underlying calls. This assumes
|
||||
// those calls need to be Qcx aware, so this takes that into account.
|
||||
type FeatureBaseImporter struct {
|
||||
api *featurebase.API
|
||||
}
|
||||
|
||||
func NewFeatureBaseImporter(api *featurebase.API) *FeatureBaseImporter {
|
||||
return &FeatureBaseImporter{
|
||||
api: api,
|
||||
}
|
||||
}
|
||||
|
||||
func (fi *FeatureBaseImporter) CreateIndexKeys(ctx context.Context, index string, keys ...string) (map[string]uint64, error) {
|
||||
return fi.api.CreateIndexKeys(ctx, index, keys...)
|
||||
}
|
||||
|
||||
func (fi *FeatureBaseImporter) CreateFieldKeys(ctx context.Context, index, field string, keys ...string) (map[string]uint64, error) {
|
||||
return fi.api.CreateFieldKeys(ctx, index, field, keys...)
|
||||
}
|
||||
|
||||
func (fi *FeatureBaseImporter) Import(ctx context.Context, req *featurebase.ImportRequest, opts ...featurebase.ImportOption) error {
|
||||
qcx := fi.api.Txf().NewQcx()
|
||||
return fi.api.Import(ctx, qcx, req, opts...)
|
||||
}
|
||||
|
||||
func (fi *FeatureBaseImporter) ImportValue(ctx context.Context, req *featurebase.ImportValueRequest, opts ...featurebase.ImportOption) error {
|
||||
qcx := fi.api.Txf().NewQcx()
|
||||
return fi.api.ImportValue(ctx, qcx, req, opts...)
|
||||
}
|
||||
|
|
@ -1,78 +0,0 @@
|
|||
package queryer
|
||||
|
||||
import (
|
||||
"context"
|
||||
|
||||
featurebase "github.com/molecula/featurebase/v3"
|
||||
"github.com/molecula/featurebase/v3/dax"
|
||||
"github.com/molecula/featurebase/v3/dax/mds/controller"
|
||||
"github.com/molecula/featurebase/v3/dax/mds/schemar"
|
||||
)
|
||||
|
||||
type MDS interface {
|
||||
// Controller-related methods.
|
||||
ComputeNodes(ctx context.Context, qtid dax.QualifiedTableID, shards ...dax.ShardNum) ([]controller.ComputeNode, error)
|
||||
IngestPartition(ctx context.Context, qtid dax.QualifiedTableID, partition dax.PartitionNum) (dax.Address, error)
|
||||
IngestShard(ctx context.Context, qtid dax.QualifiedTableID, shard dax.ShardNum) (dax.Address, error)
|
||||
TranslateNodes(ctx context.Context, qtid dax.QualifiedTableID, partitions ...dax.PartitionNum) ([]controller.TranslateNode, error)
|
||||
|
||||
// Schemar-related methods.
|
||||
schemar.Schemar
|
||||
}
|
||||
|
||||
// Ensure type implements interface.
|
||||
var _ MDS = &NopMDS{}
|
||||
|
||||
// NopMDS is a no-op implementation of the MDS interface.
|
||||
type NopMDS struct {
|
||||
schemar.NopSchemar
|
||||
}
|
||||
|
||||
func NewNopMDS() *NopMDS {
|
||||
return &NopMDS{}
|
||||
}
|
||||
|
||||
func (m *NopMDS) ComputeNodes(ctx context.Context, qtid dax.QualifiedTableID, shards ...dax.ShardNum) ([]controller.ComputeNode, error) {
|
||||
return nil, nil
|
||||
}
|
||||
func (m *NopMDS) IngestPartition(ctx context.Context, qtid dax.QualifiedTableID, partition dax.PartitionNum) (dax.Address, error) {
|
||||
return "", nil
|
||||
}
|
||||
func (m *NopMDS) IngestShard(ctx context.Context, qtid dax.QualifiedTableID, shard dax.ShardNum) (dax.Address, error) {
|
||||
return "", nil
|
||||
}
|
||||
func (m *NopMDS) TranslateNodes(ctx context.Context, qtid dax.QualifiedTableID, partitions ...dax.PartitionNum) ([]controller.TranslateNode, error) {
|
||||
return nil, nil
|
||||
}
|
||||
|
||||
type Importer interface {
|
||||
CreateIndexKeys(ctx context.Context, index string, keys ...string) (map[string]uint64, error)
|
||||
CreateFieldKeys(ctx context.Context, index, field string, keys ...string) (map[string]uint64, error)
|
||||
Import(ctx context.Context, req *featurebase.ImportRequest, opts ...featurebase.ImportOption) error
|
||||
ImportValue(ctx context.Context, req *featurebase.ImportValueRequest, opts ...featurebase.ImportOption) error
|
||||
}
|
||||
|
||||
// Ensure type implements interface.
|
||||
var _ Importer = &NopImporter{}
|
||||
|
||||
// NopImporter is a no-op implementation of the Importer interface.
|
||||
type NopImporter struct{}
|
||||
|
||||
func NewNopImporter() *NopImporter {
|
||||
return &NopImporter{}
|
||||
}
|
||||
|
||||
func (n *NopImporter) CreateIndexKeys(ctx context.Context, index string, keys ...string) (map[string]uint64, error) {
|
||||
return nil, nil
|
||||
}
|
||||
|
||||
func (n *NopImporter) CreateFieldKeys(ctx context.Context, index, field string, keys ...string) (map[string]uint64, error) {
|
||||
return nil, nil
|
||||
}
|
||||
|
||||
func (n *NopImporter) Import(ctx context.Context, req *featurebase.ImportRequest, opts ...featurebase.ImportOption) error {
|
||||
return nil
|
||||
}
|
||||
func (n *NopImporter) ImportValue(ctx context.Context, req *featurebase.ImportValueRequest, opts ...featurebase.ImportOption) error {
|
||||
return nil
|
||||
}
|
||||
File diff suppressed because it is too large
Load diff
|
|
@ -6,6 +6,7 @@ import (
|
|||
"fmt"
|
||||
"net/http"
|
||||
"strings"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
featurebase "github.com/molecula/featurebase/v3"
|
||||
|
|
@ -30,9 +31,13 @@ import (
|
|||
// that the externally-facing Molecula API would proxy query requests to a pool
|
||||
// of "Queryer" nodes, which handle incoming query requests.
|
||||
type Queryer struct {
|
||||
orchestrator *orchestrator
|
||||
mu sync.RWMutex
|
||||
orchestrators map[dax.TableQualifier]*qualifiedOrchestrator
|
||||
|
||||
mds MDS
|
||||
fbClient *featurebase.InternalClient
|
||||
|
||||
noder dax.Noder
|
||||
schemar dax.Schemar
|
||||
|
||||
logger logger.Logger
|
||||
}
|
||||
|
|
@ -40,9 +45,10 @@ type Queryer struct {
|
|||
// New returns a new instance of Queryer.
|
||||
func New(cfg Config) *Queryer {
|
||||
q := &Queryer{
|
||||
mds: NewNopMDS(),
|
||||
orchestrator: nil,
|
||||
logger: logger.NopLogger,
|
||||
noder: dax.NewNopNoder(),
|
||||
schemar: dax.NewNopSchemar(),
|
||||
orchestrators: make(map[dax.TableQualifier]*qualifiedOrchestrator),
|
||||
logger: logger.NopLogger,
|
||||
}
|
||||
|
||||
if cfg.Logger != nil {
|
||||
|
|
@ -52,8 +58,63 @@ func New(cfg Config) *Queryer {
|
|||
return q
|
||||
}
|
||||
|
||||
func (q *Queryer) SetMDS(mds MDS) error {
|
||||
q.mds = mds
|
||||
// Orchestrator gets (or creates) an instance of qualifiedOrchestrator based on
|
||||
// the provided dax.TableQualifier.
|
||||
func (q *Queryer) Orchestrator(qual dax.TableQualifier) *qualifiedOrchestrator {
|
||||
// Try to get orchestrator under a read lock first.
|
||||
if orch := func() *qualifiedOrchestrator {
|
||||
q.mu.RLock()
|
||||
defer q.mu.RUnlock()
|
||||
if orch, ok := q.orchestrators[qual]; ok {
|
||||
return orch
|
||||
}
|
||||
return nil
|
||||
}(); orch != nil {
|
||||
return orch
|
||||
}
|
||||
|
||||
// Since we didn't find an orchestrator under a read lock, obtain a write
|
||||
// lock and try a read/write.
|
||||
q.mu.Lock()
|
||||
defer q.mu.Unlock()
|
||||
if orch, ok := q.orchestrators[qual]; ok {
|
||||
return orch
|
||||
}
|
||||
|
||||
sapi := newQualifiedSchemaAPI(qual, q.schemar)
|
||||
|
||||
orch := &orchestrator{
|
||||
schema: sapi,
|
||||
trans: NewMDSTranslator(q.noder, q.schemar),
|
||||
topology: &MDSTopology{noder: q.noder},
|
||||
// TODO(jaffee) using default http.Client probably bad... need to set some timeouts.
|
||||
client: q.fbClient,
|
||||
stats: stats.NopStatsClient,
|
||||
logger: q.logger,
|
||||
}
|
||||
|
||||
qorch := newQualifiedOrchestrator(orch, qual)
|
||||
q.orchestrators[qual] = qorch
|
||||
|
||||
return qorch
|
||||
}
|
||||
|
||||
func (q *Queryer) SetNoder(noder dax.Noder) error {
|
||||
q.noder = noder
|
||||
return nil
|
||||
}
|
||||
|
||||
func (q *Queryer) SetSchemar(schemar dax.Schemar) error {
|
||||
q.schemar = schemar
|
||||
return nil
|
||||
}
|
||||
|
||||
func (q *Queryer) Start() error {
|
||||
if q.noder == nil {
|
||||
return errors.New(errors.ErrUncoded, "queryer requires noder to be configured")
|
||||
} else if q.schemar == nil {
|
||||
return errors.New(errors.ErrUncoded, "queryer requires schemar to be configured")
|
||||
}
|
||||
|
||||
// fbClient is an instance of internal client. It's used in one place in the
|
||||
// orchestrator (o.client.QueryNode()), but in that case, the host is
|
||||
|
|
@ -67,26 +128,8 @@ func (q *Queryer) SetMDS(mds MDS) error {
|
|||
if err != nil {
|
||||
return errors.Wrap(err, "setting up internal client")
|
||||
}
|
||||
q.fbClient = fbClient
|
||||
|
||||
q.orchestrator = &orchestrator{
|
||||
schema: NewSchemaInfoAPI(q.mds),
|
||||
trans: NewMDSTranslator(q.mds),
|
||||
topology: &MDSTopology{mds: q.mds},
|
||||
// TODO(jaffee) using default http.Client probably bad... need to set some timeouts.
|
||||
client: fbClient,
|
||||
stats: stats.NopStatsClient,
|
||||
logger: q.logger,
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (q *Queryer) Start() error {
|
||||
if q.mds == nil {
|
||||
return errors.New(errors.ErrUncoded, "queryer requires mds to be configured")
|
||||
} else if q.orchestrator == nil {
|
||||
return errors.New(errors.ErrUncoded, "queryer requires orchestrator to be configured")
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
|
|
@ -123,13 +166,10 @@ func (q *Queryer) QuerySQL(ctx context.Context, qual dax.TableQualifier, sql str
|
|||
}
|
||||
|
||||
// SchemaAPI
|
||||
sapi := NewQualifiedSchemaAPI(qual, q.mds)
|
||||
|
||||
// Orchestrator
|
||||
orch := newQualifiedOrchestrator(q.orchestrator, qual, q.mds)
|
||||
sapi := newQualifiedSchemaAPI(qual, q.schemar)
|
||||
|
||||
// Importer
|
||||
imp := idkmds.NewImporter(q.mds, qual, nil)
|
||||
imp := idkmds.NewImporter(q.noder, q.schemar, 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
|
||||
|
|
@ -138,7 +178,7 @@ func (q *Queryer) QuerySQL(ctx context.Context, qual dax.TableQualifier, sql str
|
|||
|
||||
systemLayer := systemlayer.NewSystemLayer()
|
||||
|
||||
pl := planner.NewExecutionPlanner(orch, sapi, sysapi, systemLayer, imp, q.orchestrator.logger, sql)
|
||||
pl := planner.NewExecutionPlanner(q.Orchestrator(qual), sapi, sysapi, systemLayer, imp, q.logger, sql)
|
||||
|
||||
planOp, err := pl.CompilePlan(ctx, st)
|
||||
if err != nil {
|
||||
|
|
@ -214,17 +254,12 @@ func (q *Queryer) QueryPQL(ctx context.Context, qual dax.TableQualifier, table d
|
|||
return nil, errors.Errorf("must have exactly 1 query, but got: %+v", qry.Calls)
|
||||
}
|
||||
|
||||
qtid, err := q.mds.TableID(ctx, qual, dax.TableName(table))
|
||||
qtbl, err := q.schemar.TableByName(ctx, qual, dax.TableName(table))
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "converting index to qualified table id")
|
||||
return nil, errors.Wrap(err, "converting index to qualified table")
|
||||
}
|
||||
|
||||
qtbl, err := q.mds.Table(ctx, qtid)
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "getting table for qtid")
|
||||
}
|
||||
|
||||
results, err := q.orchestrator.Execute(ctx, qtbl, qry, nil, &featurebase.ExecOptions{})
|
||||
results, err := q.Orchestrator(qual).Execute(ctx, qtbl, qry, nil, &featurebase.ExecOptions{})
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "orchestrator.Execute")
|
||||
}
|
||||
|
|
|
|||
|
|
@ -5,7 +5,6 @@ import (
|
|||
|
||||
pilosa "github.com/molecula/featurebase/v3"
|
||||
"github.com/molecula/featurebase/v3/dax"
|
||||
"github.com/molecula/featurebase/v3/dax/mds/schemar"
|
||||
"github.com/molecula/featurebase/v3/errors"
|
||||
)
|
||||
|
||||
|
|
@ -18,34 +17,28 @@ var _ pilosa.SchemaAPI = (*qualifiedSchemaAPI)(nil)
|
|||
// that lookup/conversion.
|
||||
type qualifiedSchemaAPI struct {
|
||||
qual dax.TableQualifier
|
||||
schemar schemar.Schemar
|
||||
schemar dax.Schemar
|
||||
}
|
||||
|
||||
func NewQualifiedSchemaAPI(qual dax.TableQualifier, schemar schemar.Schemar) *qualifiedSchemaAPI {
|
||||
func newQualifiedSchemaAPI(qual dax.TableQualifier, schema dax.Schemar) *qualifiedSchemaAPI {
|
||||
return &qualifiedSchemaAPI{
|
||||
qual: qual,
|
||||
schemar: schemar,
|
||||
schemar: schema,
|
||||
}
|
||||
}
|
||||
|
||||
func (s *qualifiedSchemaAPI) TableByName(ctx context.Context, tname dax.TableName) (*dax.Table, error) {
|
||||
qtid, err := s.schemar.TableID(ctx, s.qual, tname)
|
||||
qtbl, err := s.schemar.TableByName(ctx, s.qual, tname)
|
||||
if err != nil {
|
||||
return nil, errors.Wrapf(err, "getting table id: (%s) %s", s.qual, tname)
|
||||
}
|
||||
|
||||
qtbl, err := s.schemar.Table(ctx, qtid)
|
||||
if err != nil {
|
||||
return nil, errors.Wrapf(err, "getting table: %s", qtid)
|
||||
}
|
||||
|
||||
return &qtbl.Table, nil
|
||||
}
|
||||
|
||||
func (s *qualifiedSchemaAPI) TableByID(ctx context.Context, tid dax.TableID) (*dax.Table, error) {
|
||||
qtid := dax.NewQualifiedTableID(s.qual, tid)
|
||||
|
||||
qtbl, err := s.schemar.Table(ctx, qtid)
|
||||
qtbl, err := s.schemar.TableByID(ctx, qtid)
|
||||
if err != nil {
|
||||
return nil, errors.Wrapf(err, "getting table: %s", qtid)
|
||||
}
|
||||
|
|
@ -73,28 +66,28 @@ func (s *qualifiedSchemaAPI) CreateTable(ctx context.Context, tbl *dax.Table) er
|
|||
}
|
||||
|
||||
func (s *qualifiedSchemaAPI) CreateField(ctx context.Context, tname dax.TableName, fld *dax.Field) error {
|
||||
qtid, err := s.schemar.TableID(ctx, s.qual, tname)
|
||||
qtbl, err := s.schemar.TableByName(ctx, s.qual, tname)
|
||||
if err != nil {
|
||||
return errors.Wrapf(err, "getting table id: (%s) %s", s.qual, tname)
|
||||
return errors.Wrapf(err, "getting table by name: (%s) %s", s.qual, tname)
|
||||
}
|
||||
|
||||
return s.schemar.CreateField(ctx, qtid, fld)
|
||||
return s.schemar.CreateField(ctx, qtbl.QualifiedID(), fld)
|
||||
}
|
||||
|
||||
func (s *qualifiedSchemaAPI) DeleteTable(ctx context.Context, tname dax.TableName) error {
|
||||
qtid, err := s.schemar.TableID(ctx, s.qual, tname)
|
||||
qtbl, err := s.schemar.TableByName(ctx, s.qual, tname)
|
||||
if err != nil {
|
||||
return errors.Wrapf(err, "getting table id: (%s) %s", s.qual, tname)
|
||||
return errors.Wrapf(err, "getting table by name: (%s) %s", s.qual, tname)
|
||||
}
|
||||
|
||||
return s.schemar.DropTable(ctx, qtid)
|
||||
return s.schemar.DropTable(ctx, qtbl.QualifiedID())
|
||||
}
|
||||
|
||||
func (s *qualifiedSchemaAPI) DeleteField(ctx context.Context, tname dax.TableName, fname dax.FieldName) error {
|
||||
qtid, err := s.schemar.TableID(ctx, s.qual, tname)
|
||||
qtid, err := s.schemar.TableByName(ctx, s.qual, tname)
|
||||
if err != nil {
|
||||
return errors.Wrapf(err, "getting table id: (%s) %s", s.qual, tname)
|
||||
return errors.Wrapf(err, "getting table by name: (%s) %s", s.qual, tname)
|
||||
}
|
||||
|
||||
return s.schemar.DropField(ctx, qtid, fname)
|
||||
return s.schemar.DropField(ctx, qtid.Key().QualifiedTableID(), fname)
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1,79 +0,0 @@
|
|||
package queryer
|
||||
|
||||
import (
|
||||
"context"
|
||||
|
||||
pilosa "github.com/molecula/featurebase/v3"
|
||||
"github.com/molecula/featurebase/v3/dax"
|
||||
"github.com/molecula/featurebase/v3/dax/mds/schemar"
|
||||
"github.com/molecula/featurebase/v3/errors"
|
||||
)
|
||||
|
||||
// Ensure type implements interface.
|
||||
var _ pilosa.SchemaInfoAPI = (*schemaInfoAPI)(nil)
|
||||
|
||||
type schemaInfoAPI struct {
|
||||
schemar schemar.Schemar
|
||||
}
|
||||
|
||||
func NewSchemaInfoAPI(schemar schemar.Schemar) *schemaInfoAPI {
|
||||
return &schemaInfoAPI{
|
||||
schemar: schemar,
|
||||
}
|
||||
}
|
||||
|
||||
func (a *schemaInfoAPI) IndexInfo(ctx context.Context, indexName string) (*pilosa.IndexInfo, error) {
|
||||
qtid := dax.TableKey(indexName).QualifiedTableID()
|
||||
tbl, err := a.schemar.Table(ctx, qtid)
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "getting table for indexinfo")
|
||||
}
|
||||
|
||||
return daxTableToFeaturebaseIndexInfo(tbl, false)
|
||||
}
|
||||
|
||||
func (a *schemaInfoAPI) FieldInfo(ctx context.Context, indexName, fieldName string) (*pilosa.FieldInfo, error) {
|
||||
qtid := dax.TableKey(indexName).QualifiedTableID()
|
||||
tbl, err := a.schemar.Table(ctx, qtid)
|
||||
fldName := dax.FieldName(fieldName)
|
||||
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "getting table for fieldinfo")
|
||||
}
|
||||
|
||||
fld, ok := tbl.Field(dax.FieldName(fieldName))
|
||||
if !ok {
|
||||
return nil, dax.NewErrFieldDoesNotExist(fldName)
|
||||
}
|
||||
|
||||
return pilosa.FieldToFieldInfo(fld), nil
|
||||
}
|
||||
|
||||
// TODO(tlt): try to get rid of this in favor of pilosa.TableToIndexInfo.
|
||||
// daxTableToFeaturebaseIndexInfo converts a dax.Table to a
|
||||
// featurebase.IndexInfo. If useName is true, the IndexInfo.Name value will
|
||||
// be set to the qualified table name. Otherwise it will be set to the table key.
|
||||
func daxTableToFeaturebaseIndexInfo(qtbl *dax.QualifiedTable, useName bool) (*pilosa.IndexInfo, error) {
|
||||
name := string(qtbl.Key())
|
||||
if useName {
|
||||
name = string(qtbl.Name)
|
||||
}
|
||||
ii := &pilosa.IndexInfo{
|
||||
Name: name,
|
||||
CreatedAt: 0,
|
||||
Options: pilosa.IndexOptions{
|
||||
Keys: qtbl.StringKeys(),
|
||||
TrackExistence: true,
|
||||
},
|
||||
ShardWidth: pilosa.ShardWidth,
|
||||
}
|
||||
|
||||
// fields
|
||||
fields := make([]*pilosa.FieldInfo, len(qtbl.Fields))
|
||||
for i := range qtbl.Fields {
|
||||
fields[i] = pilosa.FieldToFieldInfo(qtbl.Fields[i])
|
||||
}
|
||||
ii.Fields = fields
|
||||
|
||||
return ii, nil
|
||||
}
|
||||
|
|
@ -50,6 +50,8 @@ func (q *queryerService) HTTPHandler() http.Handler {
|
|||
}
|
||||
|
||||
func (q *queryerService) SetMDS(addr dax.Address) error {
|
||||
q.queryer.SetMDS(mdsclient.New(addr, q.logger))
|
||||
mdscli := mdsclient.New(addr, q.logger)
|
||||
q.queryer.SetNoder(mdscli)
|
||||
q.queryer.SetSchemar(mdscli)
|
||||
return nil
|
||||
}
|
||||
|
|
|
|||
|
|
@ -14,15 +14,17 @@ import (
|
|||
)
|
||||
|
||||
// Ensure type implements interface.
|
||||
var _ Translator = (*MDSTranslator)(nil)
|
||||
var _ Translator = (*mdsTranslator)(nil)
|
||||
|
||||
type MDSTranslator struct {
|
||||
mds MDS
|
||||
type mdsTranslator struct {
|
||||
noder dax.Noder
|
||||
schemar dax.Schemar
|
||||
}
|
||||
|
||||
func NewMDSTranslator(mds MDS) *MDSTranslator {
|
||||
return &MDSTranslator{
|
||||
mds: mds,
|
||||
func NewMDSTranslator(noder dax.Noder, schemar dax.Schemar) *mdsTranslator {
|
||||
return &mdsTranslator{
|
||||
noder: noder,
|
||||
schemar: schemar,
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -37,11 +39,11 @@ func fbClient(address dax.Address) (*featurebase_client.Client, error) {
|
|||
)
|
||||
}
|
||||
|
||||
func (m *MDSTranslator) CreateIndexKeys(ctx context.Context, table string, keys []string) (map[string]uint64, error) {
|
||||
func (m *mdsTranslator) CreateIndexKeys(ctx context.Context, table string, keys []string) (map[string]uint64, error) {
|
||||
tkey := dax.TableKey(table)
|
||||
qtid := tkey.QualifiedTableID()
|
||||
|
||||
qtbl, err := m.mds.Table(ctx, qtid)
|
||||
qtbl, err := m.schemar.TableByID(ctx, qtid)
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "getting table")
|
||||
}
|
||||
|
|
@ -53,7 +55,7 @@ func (m *MDSTranslator) CreateIndexKeys(ctx context.Context, table string, keys
|
|||
|
||||
out := make(map[string]uint64)
|
||||
for pNum := range pMap {
|
||||
address, err := m.mds.IngestPartition(ctx, qtid, pNum)
|
||||
address, err := m.noder.IngestPartition(ctx, qtid, pNum)
|
||||
if err != nil {
|
||||
return nil, errors.Wrapf(err, "calling ingest-partition on table: %s, partition: %d", table, pNum)
|
||||
}
|
||||
|
|
@ -78,9 +80,9 @@ func (m *MDSTranslator) CreateIndexKeys(ctx context.Context, table string, keys
|
|||
return out, nil
|
||||
}
|
||||
|
||||
func (m *MDSTranslator) CreateFieldKeys(ctx context.Context, table string, field string, keys []string) (map[string]uint64, error) {
|
||||
func (m *mdsTranslator) CreateFieldKeys(ctx context.Context, table string, field string, keys []string) (map[string]uint64, error) {
|
||||
qtid := dax.TableKey(table).QualifiedTableID()
|
||||
address, err := m.mds.IngestPartition(ctx, qtid, dax.PartitionNum(0))
|
||||
address, err := m.noder.IngestPartition(ctx, qtid, dax.PartitionNum(0))
|
||||
if err != nil {
|
||||
return nil, errors.Wrapf(err, "calling ingest-partition on table: %s, partition: %d", table, dax.PartitionNum(0))
|
||||
}
|
||||
|
|
@ -96,11 +98,11 @@ func (m *MDSTranslator) CreateFieldKeys(ctx context.Context, table string, field
|
|||
return fbClient.CreateFieldKeys(fld, keys...)
|
||||
}
|
||||
|
||||
func (m *MDSTranslator) FindIndexKeys(ctx context.Context, table string, keys []string) (map[string]uint64, error) {
|
||||
func (m *mdsTranslator) FindIndexKeys(ctx context.Context, table string, keys []string) (map[string]uint64, error) {
|
||||
tkey := dax.TableKey(table)
|
||||
qtid := tkey.QualifiedTableID()
|
||||
|
||||
qtbl, err := m.mds.Table(ctx, qtid)
|
||||
qtbl, err := m.schemar.TableByID(ctx, qtid)
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "getting table")
|
||||
}
|
||||
|
|
@ -115,7 +117,7 @@ func (m *MDSTranslator) FindIndexKeys(ctx context.Context, table string, keys []
|
|||
pNums = append(pNums, k)
|
||||
}
|
||||
|
||||
translateNodes, err := m.mds.TranslateNodes(ctx, qtid, pNums...)
|
||||
translateNodes, err := m.noder.TranslateNodes(ctx, qtid, pNums...)
|
||||
if err != nil {
|
||||
return nil, errors.Wrapf(err, "getting translate nodes for partitions on table: %s", table)
|
||||
}
|
||||
|
|
@ -149,9 +151,9 @@ func (m *MDSTranslator) FindIndexKeys(ctx context.Context, table string, keys []
|
|||
return out, nil
|
||||
}
|
||||
|
||||
func (m *MDSTranslator) FindFieldKeys(ctx context.Context, table, field string, keys []string) (map[string]uint64, error) {
|
||||
func (m *mdsTranslator) FindFieldKeys(ctx context.Context, table, field string, keys []string) (map[string]uint64, error) {
|
||||
qtid := dax.TableKey(table).QualifiedTableID()
|
||||
address, err := m.mds.IngestPartition(ctx, qtid, dax.PartitionNum(0))
|
||||
address, err := m.noder.IngestPartition(ctx, qtid, dax.PartitionNum(0))
|
||||
if err != nil {
|
||||
return nil, errors.Wrapf(err, "calling ingest-partition on table: %s, partition: %d", table, dax.PartitionNum(0))
|
||||
}
|
||||
|
|
@ -167,7 +169,7 @@ func (m *MDSTranslator) FindFieldKeys(ctx context.Context, table, field string,
|
|||
return fbClient.FindFieldKeys(fld, keys...)
|
||||
}
|
||||
|
||||
func (m *MDSTranslator) TranslateIndexIDs(ctx context.Context, index string, ids []uint64) ([]string, error) {
|
||||
func (m *mdsTranslator) TranslateIndexIDs(ctx context.Context, index string, ids []uint64) ([]string, error) {
|
||||
idsByPartition := splitIDsByPartition(index, ids, 1<<20) // TODO(jaffee), don't hardcode shardwidth...need to get this from index info
|
||||
daxPartitions := make([]dax.PartitionNum, 0)
|
||||
for partition := range idsByPartition {
|
||||
|
|
@ -176,7 +178,7 @@ func (m *MDSTranslator) TranslateIndexIDs(ctx context.Context, index string, ids
|
|||
|
||||
qtid := dax.TableKey(index).QualifiedTableID()
|
||||
|
||||
nodes, err := m.mds.TranslateNodes(ctx, qtid, daxPartitions...)
|
||||
nodes, err := m.noder.TranslateNodes(ctx, qtid, daxPartitions...)
|
||||
if err != nil {
|
||||
return nil, errors.Wrapf(err, "calling translate-nodes on table: %s, partitions: %v", index, daxPartitions)
|
||||
}
|
||||
|
|
@ -210,7 +212,7 @@ func (m *MDSTranslator) TranslateIndexIDs(ctx context.Context, index string, ids
|
|||
return ret, nil
|
||||
}
|
||||
|
||||
func (m *MDSTranslator) TranslateIndexIDSet(ctx context.Context, table string, ids map[uint64]struct{}) (map[uint64]string, error) {
|
||||
func (m *mdsTranslator) TranslateIndexIDSet(ctx context.Context, table string, ids map[uint64]struct{}) (map[uint64]string, error) {
|
||||
idList := make([]uint64, 0, len(ids))
|
||||
for id := range ids {
|
||||
idList = append(idList, id)
|
||||
|
|
@ -227,7 +229,7 @@ func (m *MDSTranslator) TranslateIndexIDSet(ctx context.Context, table string, i
|
|||
}
|
||||
return ret, nil
|
||||
}
|
||||
func (m *MDSTranslator) TranslateFieldIDs(ctx context.Context, table, field string, ids map[uint64]struct{}) (map[uint64]string, error) {
|
||||
func (m *mdsTranslator) TranslateFieldIDs(ctx context.Context, table, field string, ids map[uint64]struct{}) (map[uint64]string, error) {
|
||||
idList := make([]uint64, 0, len(ids))
|
||||
for id := range ids {
|
||||
idList = append(idList, id)
|
||||
|
|
@ -244,9 +246,9 @@ func (m *MDSTranslator) TranslateFieldIDs(ctx context.Context, table, field stri
|
|||
}
|
||||
return ret, nil
|
||||
}
|
||||
func (m *MDSTranslator) TranslateFieldListIDs(ctx context.Context, index, field string, ids []uint64) ([]string, error) {
|
||||
func (m *mdsTranslator) TranslateFieldListIDs(ctx context.Context, index, field string, ids []uint64) ([]string, error) {
|
||||
qtid := dax.TableKey(index).QualifiedTableID()
|
||||
address, err := m.mds.IngestPartition(ctx, qtid, dax.PartitionNum(0))
|
||||
address, err := m.noder.IngestPartition(ctx, qtid, dax.PartitionNum(0))
|
||||
if err != nil {
|
||||
return nil, errors.Wrapf(err, "calling ingest-partition on table: %s, partition: %d", index, dax.PartitionNum(0))
|
||||
}
|
||||
|
|
|
|||
51
dax/schema.go
Normal file
51
dax/schema.go
Normal file
|
|
@ -0,0 +1,51 @@
|
|||
package dax
|
||||
|
||||
import "context"
|
||||
|
||||
// Schemar is similar to the pilosa.SchemaAPI interface, but it takes
|
||||
// TableQualifiers into account.
|
||||
type Schemar interface {
|
||||
TableByName(ctx context.Context, qual TableQualifier, tname TableName) (*QualifiedTable, error)
|
||||
TableByID(ctx context.Context, qtid QualifiedTableID) (*QualifiedTable, error)
|
||||
Tables(ctx context.Context, qual TableQualifier, tids ...TableID) ([]*QualifiedTable, error)
|
||||
|
||||
CreateTable(ctx context.Context, qtbl *QualifiedTable) error
|
||||
CreateField(ctx context.Context, qtid QualifiedTableID, fld *Field) error
|
||||
|
||||
DropTable(ctx context.Context, qtid QualifiedTableID) error
|
||||
DropField(ctx context.Context, qtid QualifiedTableID, fname FieldName) error
|
||||
}
|
||||
|
||||
//////////////////////////////////////////////
|
||||
|
||||
// Ensure type implements interface.
|
||||
var _ Schemar = &NopSchemar{}
|
||||
|
||||
// NopSchemar is a no-op implementation of the Schemar interface.
|
||||
type NopSchemar struct{}
|
||||
|
||||
func NewNopSchemar() *NopSchemar {
|
||||
return &NopSchemar{}
|
||||
}
|
||||
|
||||
func (s *NopSchemar) TableByName(context.Context, TableQualifier, TableName) (*QualifiedTable, error) {
|
||||
return nil, nil
|
||||
}
|
||||
func (s *NopSchemar) TableByID(ctx context.Context, qtid QualifiedTableID) (*QualifiedTable, error) {
|
||||
return nil, nil
|
||||
}
|
||||
func (s *NopSchemar) Tables(ctx context.Context, qual TableQualifier, tids ...TableID) ([]*QualifiedTable, error) {
|
||||
return nil, nil
|
||||
}
|
||||
func (s *NopSchemar) CreateTable(ctx context.Context, qtbl *QualifiedTable) error {
|
||||
return nil
|
||||
}
|
||||
func (s *NopSchemar) DropTable(ctx context.Context, qtid QualifiedTableID) error {
|
||||
return nil
|
||||
}
|
||||
func (s *NopSchemar) CreateField(ctx context.Context, qtid QualifiedTableID, fld *Field) error {
|
||||
return nil
|
||||
}
|
||||
func (s *NopSchemar) DropField(ctx context.Context, qtid QualifiedTableID, fld FieldName) error {
|
||||
return nil
|
||||
}
|
||||
|
|
@ -1027,7 +1027,7 @@ func (m *Main) setupClient() (*tls.Config, error) {
|
|||
m.SchemaManager = mds.NewSchemaManager(dax.Address(m.MDSAddress), qual, m.log)
|
||||
|
||||
m.NewImporterFn = func() pilosacore.Importer {
|
||||
return mds.NewImporter(mdsClient, qtbl.Qualifier(), &qtbl.Table)
|
||||
return mds.NewImporter(mdsClient, mdsClient, qtbl.Qualifier(), &qtbl.Table)
|
||||
}
|
||||
} else {
|
||||
m.SchemaManager = m.client
|
||||
|
|
|
|||
|
|
@ -58,7 +58,7 @@ func configureTestFlagsMDS(main *Main, address dax.Address, qtbl *dax.QualifiedT
|
|||
|
||||
mdsClient := mdsclient.New(dax.Address(address), logger.StderrLogger)
|
||||
main.NewImporterFn = func() pilosa.Importer {
|
||||
return mds.NewImporter(mdsClient, qtbl.TableQualifier, &qtbl.Table)
|
||||
return mds.NewImporter(mdsClient, mdsClient, qtbl.TableQualifier, &qtbl.Table)
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -18,18 +18,20 @@ var _ featurebase.Importer = &importer{}
|
|||
|
||||
// importer
|
||||
type importer struct {
|
||||
mds MDS
|
||||
noder dax.Noder
|
||||
schemar dax.Schemar
|
||||
|
||||
mu sync.Mutex
|
||||
qual dax.TableQualifier
|
||||
tbl *dax.Table
|
||||
}
|
||||
|
||||
func NewImporter(mds MDS, qual dax.TableQualifier, tbl *dax.Table) *importer {
|
||||
func NewImporter(noder dax.Noder, schemar dax.Schemar, qual dax.TableQualifier, tbl *dax.Table) *importer {
|
||||
return &importer{
|
||||
mds: mds,
|
||||
qual: qual,
|
||||
tbl: tbl,
|
||||
noder: noder,
|
||||
schemar: schemar,
|
||||
qual: qual,
|
||||
tbl: tbl,
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -75,7 +77,7 @@ func (m *importer) CreateTableKeys(ctx context.Context, tid dax.TableID, keys ..
|
|||
// 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.mds.IngestPartition(context.Background(), qtbl.QualifiedID(), partition)
|
||||
address, err := m.noder.IngestPartition(context.Background(), qtbl.QualifiedID(), partition)
|
||||
if err != nil {
|
||||
return nil, errors.Wrapf(err, "calling ingest-partition on table: %s, partition: %d", qtbl, partition)
|
||||
}
|
||||
|
|
@ -112,7 +114,7 @@ func (m *importer) CreateFieldKeys(ctx context.Context, tid dax.TableID, fname d
|
|||
// different partitionN for field translation.
|
||||
partition := dax.PartitionNum(0)
|
||||
|
||||
address, err := m.mds.IngestPartition(context.Background(), qtbl.QualifiedID(), partition)
|
||||
address, err := m.noder.IngestPartition(context.Background(), qtbl.QualifiedID(), partition)
|
||||
if err != nil {
|
||||
return nil, errors.Wrapf(err, "calling ingest-partition on table: %s, partition: %d", qtbl, partition)
|
||||
}
|
||||
|
|
@ -137,7 +139,7 @@ func (m *importer) ImportRoaringBitmap(ctx context.Context, tid dax.TableID, fld
|
|||
return errors.Wrapf(err, "getting qtbl")
|
||||
}
|
||||
|
||||
address, err := m.mds.IngestShard(context.Background(), qtbl.QualifiedID(), dax.ShardNum(shard))
|
||||
address, err := m.noder.IngestShard(context.Background(), qtbl.QualifiedID(), dax.ShardNum(shard))
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "calling ingest-shard")
|
||||
}
|
||||
|
|
@ -162,7 +164,7 @@ func (m *importer) ImportRoaringShard(ctx context.Context, tid dax.TableID, shar
|
|||
return errors.Wrapf(err, "getting qtbl")
|
||||
}
|
||||
|
||||
address, err := m.mds.IngestShard(context.Background(), qtbl.QualifiedID(), dax.ShardNum(shard))
|
||||
address, err := m.noder.IngestShard(context.Background(), qtbl.QualifiedID(), dax.ShardNum(shard))
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "calling ingest-shard")
|
||||
}
|
||||
|
|
@ -182,7 +184,7 @@ func (m *importer) EncodeImportValues(ctx context.Context, tid dax.TableID, fld
|
|||
return "", nil, errors.Wrapf(err, "getting qtbl")
|
||||
}
|
||||
|
||||
address, err := m.mds.IngestShard(context.Background(), qtbl.QualifiedID(), dax.ShardNum(shard))
|
||||
address, err := m.noder.IngestShard(context.Background(), qtbl.QualifiedID(), dax.ShardNum(shard))
|
||||
if err != nil {
|
||||
return "", nil, errors.Wrap(err, "calling ingest-shard")
|
||||
}
|
||||
|
|
@ -207,7 +209,7 @@ func (m *importer) EncodeImport(ctx context.Context, tid dax.TableID, fld *dax.F
|
|||
return "", nil, errors.Wrapf(err, "getting qtbl")
|
||||
}
|
||||
|
||||
address, err := m.mds.IngestShard(context.Background(), qtbl.QualifiedID(), dax.ShardNum(shard))
|
||||
address, err := m.noder.IngestShard(context.Background(), qtbl.QualifiedID(), dax.ShardNum(shard))
|
||||
if err != nil {
|
||||
return "", nil, errors.Wrap(err, "calling ingest-shard")
|
||||
}
|
||||
|
|
@ -232,7 +234,7 @@ func (m *importer) DoImport(ctx context.Context, tid dax.TableID, fld *dax.Field
|
|||
return errors.Wrapf(err, "getting qtbl")
|
||||
}
|
||||
|
||||
address, err := m.mds.IngestShard(context.Background(), qtbl.QualifiedID(), dax.ShardNum(shard))
|
||||
address, err := m.noder.IngestShard(context.Background(), qtbl.QualifiedID(), dax.ShardNum(shard))
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "calling ingest-shard")
|
||||
}
|
||||
|
|
@ -266,7 +268,7 @@ func (m *importer) getQtbl(ctx context.Context, tid dax.TableID) (*dax.Qualified
|
|||
|
||||
qtid := dax.NewQualifiedTableID(m.qual, tid)
|
||||
|
||||
qtbl, err := m.mds.Table(ctx, qtid)
|
||||
qtbl, err := m.schemar.TableByID(ctx, qtid)
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "getting table")
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1,21 +0,0 @@
|
|||
// Package mds contains the implementation of the SchemaManager interface.
|
||||
package mds
|
||||
|
||||
import (
|
||||
"context"
|
||||
|
||||
"github.com/molecula/featurebase/v3/dax"
|
||||
)
|
||||
|
||||
// MDS represents the MDS methods which importer uses.
|
||||
type MDS interface {
|
||||
IngestPartition(ctx context.Context, qtid dax.QualifiedTableID, partition dax.PartitionNum) (dax.Address, error)
|
||||
IngestShard(ctx context.Context, qtid dax.QualifiedTableID, shard dax.ShardNum) (dax.Address, error)
|
||||
|
||||
// Table was added so the `importer` instance (in this package) of the
|
||||
// batch.Importer interface could lookup up a table based on the name
|
||||
// provided in a method, as opposed to setting the table up front. This is
|
||||
// because in queryer, we don't know the table yet, because we haven't
|
||||
// parsed the sql yet.
|
||||
Table(ctx context.Context, qtid dax.QualifiedTableID) (*dax.QualifiedTable, error)
|
||||
}
|
||||
Loading…
Add table
Reference in a new issue