From 310584b0d8d69c9e3e5654bc2bf950109a90526a Mon Sep 17 00:00:00 2001 From: Ben Johnson Date: Mon, 27 Dec 2021 13:30:57 -0700 Subject: [PATCH] Add job & worker metrics --- executor.go | 20 ++++++++++++++++++-- server.go | 4 ++++ 2 files changed, 22 insertions(+), 2 deletions(-) diff --git a/executor.go b/executor.go index 5e03e1a72..f0adc7bc7 100644 --- a/executor.go +++ b/executor.go @@ -179,11 +179,18 @@ func newExecutor(opts ...executorOption) *executor { func (e *executor) addWorker() { e.workersWG.Add(1) - atomic.AddInt64(&e.currentWorkers, 1) + n := atomic.AddInt64(&e.currentWorkers, 1) + if e.Holder != nil { + e.Holder.Stats.Gauge("worker_total", float64(n), 0) + } + go func() { defer e.workersWG.Done() e.worker(e.work) - atomic.AddInt64(&e.currentWorkers, -1) + n := atomic.AddInt64(&e.currentWorkers, -1) + if e.Holder != nil { + e.Holder.Stats.Gauge("worker_total", float64(n), 0) + } }() } @@ -204,6 +211,14 @@ func (e *executor) Close() error { return nil } +// InitStats initializes stats counters. Must be called after Holder set. +func (e *executor) InitStats() { + if e.Holder != nil { + e.Holder.Stats.Count("job_total", 0, 0) + e.Holder.Stats.Gauge("worker_total", float64(atomic.LoadInt64(&e.currentWorkers)), 0) + } +} + // Execute executes a PQL query. func (e *executor) Execute(ctx context.Context, index string, q *pql.Query, shards []uint64, opt *execOptions) (QueryResponse, error) { span, ctx := tracing.StartSpanFromContext(ctx, "Executor.Execute") @@ -5989,6 +6004,7 @@ type job struct { func (e *executor) worker(work chan job) { for j := range work { atomic.AddUint64(&e.workCounter, 1) + e.Holder.Stats.Count("job_total", 1, 0) if j.idleHands { return } diff --git a/server.go b/server.go index 0c74a2fa5..a5d363e80 100644 --- a/server.go +++ b/server.go @@ -506,6 +506,10 @@ func NewServer(opts ...ServerOption) (*Server, error) { s.holder.schemator = s.schemator s.holder.sharder = s.sharder s.holder.serializer = s.serializer + + // Initial stats must be invoked after the executor obtains reference to the holder. + s.executor.InitStats() + return s, nil }