Cluster nodes for serverless (#2336)

* slowly making a serverless systemAPI for ClusterNodes()

* implemented some methods for fb_database_info

* fixed linting

* fixed comments
This commit is contained in:
David Kagan 2023-04-03 11:13:08 -04:00 committed by GitHub
parent b5dfb07118
commit 9e67f1dddd
No known key found for this signature in database
GPG key ID: 4AEE18F83AFDEB23
2 changed files with 63 additions and 5 deletions

View file

@ -39,6 +39,8 @@ type Queryer struct {
controller dax.Controller
systemLayer *systemlayer.SystemLayer
logger logger.Logger
}
@ -47,6 +49,7 @@ func New(cfg Config) *Queryer {
q := &Queryer{
controller: dax.NewNopController(),
orchestrators: make(map[dax.QualifiedDatabaseID]*qualifiedOrchestrator),
systemLayer: systemlayer.NewSystemLayer(),
logger: logger.NopLogger,
}
@ -198,17 +201,14 @@ func (q *Queryer) QuerySQL(ctx context.Context, qdbid dax.QualifiedDatabaseID, s
// Importer
imp := idkserverless.NewImporter(q.controller, qdbid, nil)
// TODO(tlt): We need a serverless-compatible implementation of the
// SystemAPI.
sysapi := &featurebase.NopSystemAPI{}
systemLayer := systemlayer.NewSystemLayer()
sysapi := newSystemAPI(q.controller, qdbid)
// We intentionally don't pass the sql argument here because we're working
// with an io.Reader rather than a string and it's just not necessary to
// send it as a string to this method. Also, what happens if the sql is a
// large BULK INSERT?
pl := planner.NewExecutionPlanner(q.Orchestrator(qdbid), sapi, sysapi, systemLayer, imp, q.logger, "")
pl := planner.NewExecutionPlanner(q.Orchestrator(qdbid), sapi, sysapi, q.systemLayer, imp, q.logger, "")
planOp, err := pl.CompilePlan(ctx, st)
if err != nil {

58
dax/queryer/system_api.go Normal file
View file

@ -0,0 +1,58 @@
package queryer
import (
"context"
featurebase "github.com/featurebasedb/featurebase/v3"
"github.com/featurebasedb/featurebase/v3/dax"
)
// systemAPI is an implementation of the systemAPI.
type systemAPI struct {
featurebase.NopSystemAPI
controller dax.Controller
qdbid dax.QualifiedDatabaseID
}
func newSystemAPI(c dax.Controller, qdbid dax.QualifiedDatabaseID) *systemAPI {
return &systemAPI{
controller: c,
qdbid: qdbid,
}
}
// ClusterNodes returns a list of featurebase.ClusterNodes
// with length of the minimum number of workers
func (s *systemAPI) ClusterNodes() []featurebase.ClusterNode {
ctx := context.Background()
qdb, err := s.controller.DatabaseByID(ctx, s.qdbid)
if err != nil {
return []featurebase.ClusterNode{}
}
out := make([]featurebase.ClusterNode, qdb.Options.WorkersMin)
return out
}
func (s *systemAPI) PlatformDescription() string {
return "Serverless"
}
func (s *systemAPI) ClusterName() string {
return "Serverless"
}
func (s *systemAPI) ClusterNodeCount() int {
ctx := context.Background()
qdb, err := s.controller.DatabaseByID(ctx, s.qdbid)
if err != nil {
return -1
}
return qdb.Options.WorkersMin
}
func (s *systemAPI) ClusterState() string {
return "NORMAL"
}