Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
122 changes: 122 additions & 0 deletions fracmanager/search_stats.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,122 @@
package fracmanager

import (
"fmt"
"time"

"go.uber.org/zap"

"github.com/ozontech/seq-db/frac"
"github.com/ozontech/seq-db/querytracer"
"github.com/ozontech/seq-db/seq"
)

const maxFracsSlowSearchLog = 10

type fracTimings struct {
semaphoreWait time.Duration
elapsed []time.Duration
}

type searchStats struct {
iterations int

total time.Duration
semaphoreWait time.Duration

min time.Duration
fastest string

max time.Duration
slowest string

matched int
unmatched int
matchedNames []string
unmatchedNames []string
}

func (s *searchStats) merge(fracs []frac.Fraction, qprs []*seq.QPR, timings fracTimings) {
s.iterations++
s.semaphoreWait += timings.semaphoreWait

for i, elapsed := range timings.elapsed {
if s.fracs() == 0 || elapsed < s.min {
s.min = elapsed
s.fastest = fracs[i].Info().Name()
}

if elapsed > s.max {
s.max = elapsed
s.slowest = fracs[i].Info().Name()
}

s.total += elapsed

if qprs[i].Empty() {
s.unmatched++
if len(s.unmatchedNames) < maxFracsSlowSearchLog {
s.unmatchedNames = append(s.unmatchedNames, fracs[i].Info().Name())
}
continue
}

s.matched++
if len(s.matchedNames) < maxFracsSlowSearchLog {
s.matchedNames = append(s.matchedNames, fracs[i].Info().Name())
}
}
}

func (s *searchStats) fracs() int {
return s.matched + s.unmatched
}

func (s *searchStats) mean() time.Duration {
if s.fracs() == 0 {
return 0
}
return s.total / time.Duration(s.fracs())
}

func (s *searchStats) addToTracer(tr *querytracer.Tracer) {
if !tr.Enabled() {
return
}

tr.AddChildWithSpan(&querytracer.Span{
Message: "waiting on the search worker semaphore",
Duration: s.semaphoreWait,
})

tr.AddChildWithSpan(&querytracer.Span{
Message: fmt.Sprintf(
"fraction search time: %d fractions (%d matched) in %d iterations",
s.fracs(), s.matched, s.iterations,
),
Duration: s.total,
Children: []*querytracer.Span{
{Message: "mean", Duration: s.mean()},
{Message: fmt.Sprintf("min (%s)", s.fastest), Duration: s.min},
{Message: fmt.Sprintf("max (%s)", s.slowest), Duration: s.max},
},
})
}

func (s *searchStats) zapFields() []zap.Field {
return []zap.Field{
zap.Int("search_iterations", s.iterations),
zap.Int("fracs_searched", s.fracs()),
zap.Int("fracs_matched", s.matched),
zap.Strings("fracs_matched_sample", s.matchedNames),
zap.Int("fracs_unmatched", s.unmatched),
zap.Strings("fracs_unmatched_sample", s.unmatchedNames),
zap.Int64("semaphore_wait_ms", s.semaphoreWait.Milliseconds()),
zap.Int64("frac_search_total_ms", s.total.Milliseconds()),
zap.Int64("frac_search_min_ms", s.min.Milliseconds()),
zap.String("frac_search_min_frac", s.fastest),
zap.Int64("frac_search_mean_ms", s.mean().Milliseconds()),
zap.Int64("frac_search_max_ms", s.max.Milliseconds()),
zap.String("frac_search_max_frac", s.slowest),
}
}
149 changes: 52 additions & 97 deletions fracmanager/searcher.go
Original file line number Diff line number Diff line change
Expand Up @@ -19,10 +19,6 @@ import (
"github.com/ozontech/seq-db/util"
)

const (
maxFracsSlowSearchLog = 10
)

type SearcherCfg struct {
MaxFractionHits int // the maximum number of fractions used in the search
FractionsPerIteration int
Expand Down Expand Up @@ -52,7 +48,6 @@ func (s *Searcher) SearchDocs(ctx context.Context, fracs []frac.Fraction, params
return nil, err
}

subSearchesCnt := 0
origLimit := params.Limit
scanAll := params.IsScanAllRequest()

Expand All @@ -66,85 +61,53 @@ func (s *Searcher) SearchDocs(ctx context.Context, fracs []frac.Fraction, params
fracsChunkSize = len(remainingFracs)
}

var totalSearchTimeNanos int64
var totalWaitTimeNanos int64
var totalMemUsage int
totalFracsFound := 0
totalFracsSkipped := 0
var fracsFound []string
var fracsSkipped []string
var (
stats searchStats
qprMemUsage int
)

for len(remainingFracs) > 0 && (scanAll || params.Limit > 0) {
chunk := remainingFracs.Shift(fracsChunkSize)

subQPRs, searchTimeNanos, waitTimeNanos, err := s.searchDocsAsync(ctx, chunk, params)
subQPRs, timings, err := s.searchDocsAsync(ctx, chunk, params)
if err != nil {
return nil, err
}

for i, qpr := range subQPRs {
if !qpr.Empty() {
totalFracsFound++
if len(fracsFound) < maxFracsSlowSearchLog {
fracsFound = append(fracsFound, chunk[i].Info().Name())
}
} else {
totalFracsSkipped++
if len(fracsSkipped) < maxFracsSlowSearchLog {
fracsSkipped = append(fracsSkipped, chunk[i].Info().Name())
}
}
}

totalSearchTimeNanos += searchTimeNanos
totalWaitTimeNanos += waitTimeNanos
stats.merge(chunk, subQPRs, timings)

seq.MergeQPRs(total, subQPRs, origLimit, seq.MillisToMID(params.HistInterval), params.Order)

totalMemUsage = total.MemUsage()

if s.cfg.MaxQprMemory > 0 && totalMemUsage > s.cfg.MaxQprMemory {
return nil, fmt.Errorf("%w: used %d bytes, limit %d", consts.ErrMemoryLimitExceeded, total.MemUsage(), s.cfg.MaxQprMemory)
qprMemUsage = total.MemUsage()
if s.cfg.MaxQprMemory > 0 && qprMemUsage > s.cfg.MaxQprMemory {
return nil, fmt.Errorf(
"%w: used %d bytes, limit %d",
consts.ErrMemoryLimitExceeded, qprMemUsage, s.cfg.MaxQprMemory,
)
}

// reduce the limit on the number of ensured docs in response
params.Limit = origLimit - calcEnsuredIDsCount(total.IDs, remainingFracs, params.Order)

subSearchesCnt++
}

if tr.Enabled() {
searchSpan := &querytracer.Span{
Message: "search iteratively (cpu time)",
Duration: time.Duration(totalSearchTimeNanos),
}
tr.AddChildWithSpan(searchSpan)
waitSpan := &querytracer.Span{
Message: "waiting goroutines (all cores)",
Duration: time.Duration(totalWaitTimeNanos),
}
tr.AddChildWithSpan(waitSpan)
}

searchSubSearches.Observe(float64(subSearchesCnt))
stats.addToTracer(tr)
searchSubSearches.Observe(float64(stats.iterations))

took := time.Since(start)
if s.cfg.SlowLogThreshold != 0 && took >= s.cfg.SlowLogThreshold {
fields := []zap.Field{
zap.Int64("took_ms", took.Milliseconds()),
zap.Object("params", params),
zap.Int("total_fracs_found", totalFracsFound),
zap.Strings("fracs_found", fracsFound),
zap.Int("total_fracs_skipped", totalFracsSkipped),
zap.Strings("fracs_skipped", fracsSkipped),
util.ZapUint64AsSizeStr("qpr_size", uint64(totalMemUsage)),
zap.Uint64("total", total.Total),
zap.Int64("took_ms", took.Milliseconds()),
util.ZapUint64AsSizeStr("qpr_size", uint64(qprMemUsage)),
}

logger.Warn("slow search", fields...)
logger.Warn(
"slow search",
append(fields, stats.zapFields()...)...,
)
}
return total, nil

return total, nil
}

func (s *Searcher) prepareFracs(fracs List, params processor.SearchParams) (List, error) {
Expand Down Expand Up @@ -183,64 +146,56 @@ func calcEnsuredIDsCount(ids seq.IDSources, remainingFracs List, order seq.DocsO
return sort.Search(len(ids), func(i int) bool { return ids[i].ID.MID <= nextFracInfo.To })
}

func (s *Searcher) searchDocsAsync(ctx context.Context, fracs []frac.Fraction, params processor.SearchParams) ([]*seq.QPR, int64, int64, error) {
func (s *Searcher) searchDocsAsync(
ctx context.Context,
fracs []frac.Fraction,
params processor.SearchParams,
) ([]*seq.QPR, fracTimings, error) {
ctx, cancel := context.WithCancel(ctx)
defer cancel()

var err error
var (
err error
once sync.Once
wg sync.WaitGroup
)

once := sync.Once{}
wg := sync.WaitGroup{}
qprs := make([]*seq.QPR, len(fracs))

wgDoneNanos := make([]int64, len(fracs))
searchElapsedNanos := make([]int64, len(fracs))
timings := fracTimings{elapsed: make([]time.Duration, len(fracs))}

loop:
for i, frac := range fracs {
for i, f := range fracs {
acquireStart := time.Now()

select {
case <-ctx.Done():
once.Do(func() { err = ctx.Err() })
break loop
case s.sem <- struct{}{}: // acquire semaphore
wg.Add(1)
go func() {
searchStart := time.Now()
var fracErr error
if qprs[i], fracErr = s.fracSearch(ctx, params, frac); fracErr != nil {
once.Do(func() {
err = fracErr
cancel()
})
}
searchElapsedNanos[i] = time.Since(searchStart).Nanoseconds()

<-s.sem // release semaphore

wgDoneNanos[i] = time.Now().UnixNano()
wg.Done()
}()
}
}

wg.Wait()
waitEndTime := time.Now().UnixNano()
timings.semaphoreWait += time.Since(acquireStart)

totalSearchTimeNanos := int64(0)
totalWaitTimeNanos := int64(0)
wg.Go(func() {
defer func() { <-s.sem }() // release semaphore

for i := range fracs {
if wgDoneNanos[i] != 0 {
totalWaitTimeNanos += waitEndTime - wgDoneNanos[i]
}
totalSearchTimeNanos += searchElapsedNanos[i]
}
searchStart := time.Now()
qpr, fracErr := s.fracSearch(ctx, params, f)
timings.elapsed[i] = time.Since(searchStart)

if err != nil {
return nil, 0, 0, err
if fracErr != nil {
once.Do(func() {
err = fracErr
cancel()
})
return
}
qprs[i] = qpr
})
}

return qprs, totalSearchTimeNanos, totalWaitTimeNanos, nil
wg.Wait()
return qprs, timings, err
}

func (s *Searcher) fracSearch(ctx context.Context, params processor.SearchParams, f frac.Fraction) (_ *seq.QPR, err error) {
Expand Down
Loading