Skip to content

Add cluster namespace fanout heuristics supporting queries greater than retention - #908

Merged
robskillington merged 6 commits into
masterfrom
r/local-query-cluster-fanout-heuristics
Oct 5, 2018
Merged

Add cluster namespace fanout heuristics supporting queries greater than retention#908
robskillington merged 6 commits into
masterfrom
r/local-query-cluster-fanout-heuristics

Conversation

@robskillington

@robskillington robskillington commented Sep 16, 2018

Copy link
Copy Markdown
Collaborator

I need to add tests, however I want to get buy-in that this is meets our current needs before moving to writing extensive tests.

This fixes #866.

Comment thread src/query/storage/local/storage.go Outdated
existing.attrs.Resolution <= attrs.Resolution
existsBetter = longerRetention || sameRetentionAndSameOrMoreGranularResolution
default:
panic(fmt.Sprintf("unknown query fanout type: %d", r.queryFanoutType))

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

hm could you error instead of panic-ing here?

// cluster that can completely fulfill this range and then prefer the
// highest resolution (most fine grained) results.
// This needs to be optimized, however this is a start.
fanout, namespaces, err := s.resolveClusterNamespacesForQuery(query.Start, query.End)

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

couple thoughts:

  • maybe it'd be good to add a way for users to see what data source was hit for each series
  • do you want to add tests for the various combinations - just ensuring we override to the right storage based on the heuristics in code

@codecov

codecov Bot commented Oct 4, 2018

Copy link
Copy Markdown

Codecov Report

Merging #908 into master will increase coverage by 0.02%.
The diff coverage is 76.02%.

Impacted file tree graph

@@            Coverage Diff             @@
##           master     #908      +/-   ##
==========================================
+ Coverage   76.72%   76.75%   +0.02%     
==========================================
  Files         436      436              
  Lines       37002    37086      +84     
==========================================
+ Hits        28391    28464      +73     
- Misses       6585     6590       +5     
- Partials     2026     2032       +6
Flag Coverage Δ
#dbnode 81.35% <100%> (+0.04%) ⬆️
#m3em 66.39% <ø> (ø) ⬆️
#m3ninx 75.33% <ø> (+0.07%) ⬆️
#m3nsch 51.19% <ø> (ø) ⬆️
#query 63.93% <75.69%> (+0.06%) ⬆️
#x 84.72% <ø> (ø) ⬆️

Continue to review full report at Codecov.

Legend - Click here to learn more
Δ = absolute <relative> (impact), ø = not affected, ? = missing data
Powered by Codecov. Last update 94784bd...1da3723. Read the comment docs.

// TODO: Consider combining readWorkerPool and writeWorkerPool
func NewStorage(clusters Clusters, workerPool pool.ObjectPool, writeWorkerPool xsync.PooledWorkerPool) Storage {
return &m3storage{clusters: clusters, readWorkerPool: workerPool, writeWorkerPool: writeWorkerPool}
func NewStorage(

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

i thought artem had removed the TODO above in another PR. mind rebasing on master once?

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Sure thing.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I've rebased, its still there.

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

sigh must have been skipped in the last review, mind deleting it please.

r.seenFirstAttrs = attrs
}
r.seenIters = append(r.seenIters, iterators)
if !r.err.Empty() {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

should this be moved up to the first thing after defering the unlock?

@robskillington robskillington Oct 5, 2018

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

No this still needs to be here, otherwise we won't free/close the iters for successful requests that come back.

I'm adding a note about it so that it's obvious that's why this is below.

iter: iter,
})
if len(r.seenIters) == 2 {
// need to build backfill the dedupe map from the first result first

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

maybe remove the word build

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Sure thing.

return r.finalResult, nil
}

func (r *multiResult) Add(

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

super nit: maybe call the iterators arg newIterators, I got a little lost between seenIters and iterators the first time

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Sure thing.

serieses[idx] = multiResultSeries{}
var existsBetter bool
switch r.fanout {
case namespacesCoverAllRetention:

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

is this more like namespacesCoverQueryDuration?

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Yeah I can rename this. Probably to namespaceCoversQueryRange.

case namespacesCoverPartialRetention:
// Already exists and either has longer retention, or the same retention
// and result we are adding is not as precise
longerRetention := existing.attrs.Retention > attrs.Retention

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

nit: existsLongerRetention

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Sure thing.

for idx := range serieses {
serieses[idx] = multiResultSeries{}
var existsBetter bool
switch r.fanout {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

nit: I feel like this logic would be easier to follow if it was framed from the perspective of determining whether the incoming attributes were better as opposed to vice versa, but that might just be personal preference

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Hm, I'm kind of preferring to keep the existing one if its there, so that's why I preferably found it better to frame it that way.

@richardartoul

Copy link
Copy Markdown
Contributor

One thing I'm realizing now is that there is no concept of how old a namespace is. So say you start with:

namespace1: 5 day retention, 1m resolution

Then you decide you want higher resolution so you start dual writing to:

namespace1: 5 day retention, 1m resolution
namespace2: 5 day retention, 30s resolution

Your queries will immediately begin going to namespace 2 even though it doesn't really have any data.

Fixing that is probably out of scope for this P.R but maybe we can track in an issue or something, although I'm not sure what the solution is other than configuring "cutoffs" or something

@robskillington

Copy link
Copy Markdown
Collaborator Author

@richardartoul yeah, I'm explicitly not tackling that.

We can add a readable/writeable flag later to the namespaces - and then manually flip it on once the namespace is filled up.

Then eventually we'll get to the all singing, all dancing, dynamic and automated version of this all.

r.seenIters = nil

if r.finalResult != nil {
// NB(r): Since all the series iterators in the final result are held onto

@prateek prateek Oct 5, 2018

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Mind adding a unit test w/ mocks to ensure we free the underlying iterators exactly once assuming the API is used correctly. I think it's 100% right the way it is but the test will be good for when someone comes and attempts to refactors this.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Sure thing.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

As discussed, let's revisit this in a followup change.

// need to backfill the dedupe map from the first result first
existing := r.seenIters[0]
r.dedupeMap = make(map[string]multiResultSeries, existing.Len())
for _, iter := range existing.Iters() {

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

nit: could you refactor this loop into func (r *multiResult) addOrUpdateMap(attrs, iters) - would be re-usable here and in dedupe()

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Sure thing.

Comment thread src/query/storage/m3/storage.go Outdated
start time.Time,
end time.Time,
) (queryFanoutType, ClusterNamespaces, error) {
now := time.Now()

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

maybe just set a nowFn on *m3Storage and initialize it to time.Now for now, so at least this is slightly easier to work with if someone needs to control it in a test later

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Sure thing.


unaggregated := s.clusters.UnaggregatedClusterNamespace()
unaggregatedRetention := unaggregated.Options().Attributes().Retention
unaggregatedStart := now.Add(-1 * unaggregatedRetention)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Do you not need to truncate to a blockSize here? I guess thats just kind of an implementation detail of M3DB and not super relevant here..

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Not here thankfully.

Comment thread src/query/storage/m3/storage.go Outdated
unaggregated := s.clusters.UnaggregatedClusterNamespace()
unaggregatedRetention := unaggregated.Options().Attributes().Retention
unaggregatedStart := now.Add(-1 * unaggregatedRetention)
if !unaggregatedStart.After(start) {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Is this the same as unaggregatedStart.Before(start)? If so thats easier to reason about to me

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

oh I see you would have to add || .Equal() as well but still seems nicer to me

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Hm, yeah I did it so you don't have to call twice. It's not a hot path though so I'll refactor to this.

Comment thread src/query/storage/m3/storage.go Outdated
return namespacesCoverAllRetention, ClusterNamespaces{unaggregated}, nil
}

// First determine if any aggregated clusters that can span the whole query

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

"First determine if any aggregated clusters span the whole query range, if so..."

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Good catch, thanks.

case namespaceCoversPartialQueryRange:
// Already exists and either has longer retention, or the same retention
// and result we are adding is not as precise
existsLongerRetention := existing.attrs.Retention > attrs.Retention

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

All this is to avoid consolidation... :sigh:

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Yup =]


downsampleOpts, err := opts.DownsampleOptions()
if err != nil || !downsampleOpts.All {
// Cluster does not contain all data, include as part of fan out

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Confused about the downsampling thing, how is this different than saying something is unaggregated?

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Basically it's aggregated, but data will only go to this cluster if rules are setup to do so.

@richardartoul richardartoul left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

LGTM with the nits

@robskillington
robskillington merged commit c72ee5b into master Oct 5, 2018
@robskillington
robskillington deleted the r/local-query-cluster-fanout-heuristics branch October 5, 2018 20:00
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

PromQL query with longer time range causes internal server error

3 participants