From 60d534c505f56a7ac5e4349dda933d74cb3c59f1 Mon Sep 17 00:00:00 2001 From: Ben Johnson Date: Mon, 2 Aug 2021 15:27:05 -0600 Subject: [PATCH] Limit translation memory & add max query memory config --- ctl/server.go | 1 + executor.go | 53 +++++++++++++++++++++++++++++++++++++----------- server.go | 17 +++++++++++++++- server/config.go | 3 +++ server/server.go | 3 ++- 5 files changed, 63 insertions(+), 14 deletions(-) diff --git a/ctl/server.go b/ctl/server.go index b4f8b8029..2660a8b37 100644 --- a/ctl/server.go +++ b/ctl/server.go @@ -39,6 +39,7 @@ func BuildServerFlags(cmd *cobra.Command, srv *server.Command) { flags.Uint64Var(&srv.Config.MaxFileCount, "max-file-count", srv.Config.MaxFileCount, "Soft limit on the maximum number of fragment files FeatureBase keeps open simultaneously.") flags.DurationVar((*time.Duration)(&srv.Config.LongQueryTime), "long-query-time", time.Duration(srv.Config.LongQueryTime), "Duration that will trigger log and stat messages for slow queries. Zero to disable.") flags.IntVar(&srv.Config.QueryHistoryLength, "query-history-length", srv.Config.QueryHistoryLength, "Number of queries to remember in history.") + flags.Int64Var(&srv.Config.MaxQueryMemory, "max-query-memory", srv.Config.MaxQueryMemory, "Maximum memory allowed per Extract() or SELECT query.") // TLS SetTLSConfig(flags, "", &srv.Config.TLS.CertificatePath, &srv.Config.TLS.CertificateKeyPath, &srv.Config.TLS.CACertPath, &srv.Config.TLS.SkipVerify, &srv.Config.TLS.EnableClientVerification) diff --git a/executor.go b/executor.go index b0d1959d4..ec071f34f 100644 --- a/executor.go +++ b/executor.go @@ -205,6 +205,10 @@ func (e *executor) Execute(ctx context.Context, index string, q *pql.Query, shar if opt == nil { opt = &execOptions{} } + // Default maximum memory, if not passed in. + if opt.MaxMemory == 0 { + opt.MaxMemory = e.maxMemory + } if opt.Profile { var prof tracing.ProfiledSpan @@ -235,7 +239,7 @@ func (e *executor) Execute(ctx context.Context, index string, q *pql.Query, shar // No need to translate a remote call. if !opt.Remote { // only translateResults if this local node is the final destination. only string/column keys. - if err := e.translateResults(ctx, index, idx, q.Calls, results); err != nil { + if err := e.translateResults(ctx, index, idx, q.Calls, results, opt.MaxMemory); err != nil { if errors.Cause(err) == ErrTranslatingKeyNotFound { // No error - return empty result resp.Results = make([]interface{}, len(q.Calls)) @@ -3956,7 +3960,7 @@ func (e *executor) executeExternalLookup(ctx context.Context, qcx *Qcx, index st } qr := []interface{}{rawArg} - err = e.translateResults(ctx, index, idx, c.Children, qr) + err = e.translateResults(ctx, index, idx, c.Children, qr, e.maxMemory) if err != nil { return ExtractedTable{}, errors.Wrap(err, "translating query result") } @@ -4163,11 +4167,6 @@ func (e *executor) executeExternalLookup(ctx context.Context, qcx *Qcx, index st } func (e *executor) executeExtract(ctx context.Context, qcx *Qcx, index string, c *pql.Call, shards []uint64, opt *execOptions) (ExtractedIDMatrix, error) { - // Defualt maximum memory, if not passed in. - if opt.MaxMemory == 0 { - opt.MaxMemory = e.maxMemory - } - // Extract the column filter call. if len(c.Children) < 1 { return ExtractedIDMatrix{}, errors.New("missing column filter in Extract") @@ -5852,12 +5851,13 @@ func (e *executor) mapper(ctx context.Context, eg *errgroup.Group, ch chan mapRe // calcResultMemory recursively computes the total memory used by v. func calcResultMemory(v interface{}) (n int64) { switch v := v.(type) { - case string: case ExtractedIDColumn: n += 8 // ColumnID for _, row := range v.Rows { n += 24 + int64(len(row)*8) // slice header + data } + return n + case ExtractedIDMatrix: n += 24 // slice size for _, field := range v.Fields { @@ -5868,8 +5868,34 @@ func calcResultMemory(v interface{}) (n int64) { for _, col := range v.Columns { n += calcResultMemory(col) } + return n + + case ExtractedTableColumn: + n += 8 + 16 + int64(len(v.Column.Key)) + 8 // KeyOrID + for _, row := range v.Rows { + n += 8 + calcResultMemory(row) // ptr + value size + } + return n + + case string: + return 16 + int64(len(v)) + case bool, int64, uint64: + return 8 + case []string: + n += 24 // slice header + for i := range v { + n += 16 + int64(len(v[i])) + } + return n + case []uint64: + return 24 + int64(8*len(v)) // slice header + data size + case pql.Decimal: + return 16 + case time.Time: + return 24 + default: + return n } - return n } type job struct { @@ -6617,7 +6643,7 @@ func (e *executor) callZero(c *pql.Call) *pql.Call { } } -func (e *executor) translateResults(ctx context.Context, index string, idx *Index, calls []*pql.Call, results []interface{}) (err error) { +func (e *executor) translateResults(ctx context.Context, index string, idx *Index, calls []*pql.Call, results []interface{}, memoryAvailable int64) (err error) { span, _ := tracing.StartSpanFromContext(ctx, "Executor.translateResults") defer span.Finish() @@ -6636,7 +6662,7 @@ func (e *executor) translateResults(ctx context.Context, index string, idx *Inde } for i := range results { - results[i], err = e.translateResult(ctx, index, idx, calls[i], results[i], idMap) + results[i], err = e.translateResult(ctx, index, idx, calls[i], results[i], idMap, &memoryAvailable) if err != nil { return err } @@ -6758,7 +6784,7 @@ func (e *executor) preTranslateMatrixSet(mat ExtractedIDMatrix, fieldIdx uint, f return e.Cluster.translateFieldIDs(field, ids) } -func (e *executor) translateResult(ctx context.Context, index string, idx *Index, call *pql.Call, result interface{}, idSet map[uint64]string) (_ interface{}, err error) { +func (e *executor) translateResult(ctx context.Context, index string, idx *Index, call *pql.Call, result interface{}, idSet map[uint64]string, memoryAvailable *int64) (_ interface{}, err error) { switch result := result.(type) { case *Row: rowIdx, rowField, strategy, err := e.howToTranslate(idx, result) @@ -7198,6 +7224,9 @@ func (e *executor) translateResult(ctx context.Context, index string, idx *Index Column: colTrans, Rows: data, } + if *memoryAvailable -= calcResultMemory(cols[i]); *memoryAvailable < 0 { + return nil, fmt.Errorf("table exceeds available memory") + } } return ExtractedTable{ diff --git a/server.go b/server.go index 233c61b42..b3b0b5128 100644 --- a/server.go +++ b/server.go @@ -91,6 +91,7 @@ type Server struct { // nolint: maligned confirmDownSleep time.Duration confirmDownRetries int syncer holderSyncer + maxQueryMemory int64 translationSyncer TranslationSyncer resetTranslationSyncCh chan struct{} @@ -368,6 +369,14 @@ func OptServerQueryHistoryLength(length int) ServerOption { } } +// OptServerMaxQueryMemory sets the memory used per Extract() and SELECT query. +func OptServerMaxQueryMemory(v int64) ServerOption { + return func(s *Server) error { + s.maxQueryMemory = v + return nil + } +} + // OptServerDisCo is a functional option on Server // used to set the Distributed Consensus implementation. func OptServerDisCo(disCo disco.DisCo, @@ -454,10 +463,16 @@ func NewServer(opts ...ServerOption) (*Server, error) { return nil, errors.Wrap(err, "mem total") } + // Default memory to 20% of total. + maxQueryMemory := s.maxQueryMemory + if maxQueryMemory == 0 { + maxQueryMemory = int64(float64(memTotal) * .20) + } + // set up executor after server opts have been processed executorOpts := []executorOption{ optExecutorInternalQueryClient(s.defaultClient), - optExecutorMaxMemory(int64(float64(memTotal) * .50)), + optExecutorMaxMemory(maxQueryMemory), } if s.executorPoolSize > 0 { executorOpts = append(executorOpts, optExecutorWorkerPoolSize(s.executorPoolSize)) diff --git a/server/config.go b/server/config.go index bb5b412c5..b96a5be1f 100644 --- a/server/config.go +++ b/server/config.go @@ -130,6 +130,9 @@ type Config struct { // don't exhaust the goroutine limit. ImportWorkerPoolSize int `toml:"-"` + // Limits the total amount of memory to be used by Extract() & SELECT queries. + MaxQueryMemory int64 `toml:"max-query-memory"` + Cluster struct { ReplicaN int `toml:"replicas"` Name string `toml:"name"` diff --git a/server/server.go b/server/server.go index 967dbd562..dd5488561 100644 --- a/server/server.go +++ b/server/server.go @@ -40,7 +40,6 @@ import ( "golang.org/x/sync/errgroup" - "github.com/pelletier/go-toml" "github.com/molecula/featurebase/v2" "github.com/molecula/featurebase/v2/boltdb" "github.com/molecula/featurebase/v2/encoding/proto" @@ -56,6 +55,7 @@ import ( "github.com/molecula/featurebase/v2/statsd" "github.com/molecula/featurebase/v2/syswrap" "github.com/molecula/featurebase/v2/testhook" + "github.com/pelletier/go-toml" "github.com/pkg/errors" ) @@ -493,6 +493,7 @@ func (m *Command) SetupServer() error { pilosa.OptServerStorageConfig(m.Config.Storage), pilosa.OptServerRowcacheOn(m.Config.RowcacheOn), pilosa.OptServerRBFConfig(m.Config.RBFConfig), + pilosa.OptServerMaxQueryMemory(m.Config.MaxQueryMemory), pilosa.OptServerQueryHistoryLength(m.Config.QueryHistoryLength), discoOpt, }