Add job & worker metrics

This commit is contained in:
Ben Johnson 2021-12-27 13:30:57 -07:00
parent 145f65ab0e
commit 310584b0d8
2 changed files with 22 additions and 2 deletions

View file

@ -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
}

View file

@ -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
}