Repository navigation
Feature: flow metrics integration - #14518
Conversation
* metrics: eliminate race condition when registering metrics Ensure our fast-lookup and store tables cannot diverge in a race condition by wrapping mutation of both in a single mutex and appropriately handle another thread winning the race to the lock by using the value that it persisted instead of writing our own. * metrics: guard against intermediate namespace conflicts - ensures our safeguard that prevents using an existing metric as a namespace is applied to _intermediate_ nodes, not just the tail-node, eliminating a potential crash when sending `fetch_or_store` to a metric object that is not expected to respond to `fetch_or_store`. - uses the atomic `Concurrent::Map#compute_if_absent` instead of the non-atomic `Concurrent::Map#fetch_or_store`, which is prone to last-write-wins during contention (as-written, this method is only executed under lock and not subject to contention) - uses `Enumerable#reduce` to eliminate the need for recursion * flow: introduce auto-advancing UptimeMetric * flow: introduce FlowMetric with minimal current/lifetime rates * flow: initialize pipeline metrics at pipeline start
* Controller and service layer implementation for flow metrics. * Add flow metrics to unit test and benchmark cli definitions.
yaauie
left a comment
There was a problem hiding this comment.
First pass: concurrency could be made more clear as worker_concurrency, and backpressure similarly reads more clearly as queue_backpressure.
Rename `concurrency` to `worker_concurrency ` and `backpressure` to `queue_backpressure` to provide proper scope naming. Co-authored-by: Ry Biesemeyer <yaauie@users.noreply.github.com>
the collector is absent when the pipeline is run in test with a NullMetricExt, or when the pipeline is explicitly configured to not collect metrics using `metric.collect: false`.
* Unit tests and integration tests added for flow metrics. * Node stat spec and pipeline spec metric updates. * Metric keys statically imported, implicit error expectation added in metric spec. * Fix node status API spec after renaming flow metrics. * Removing flow metric from PipelinesInfo DS (used in peridoci metric snapshot), integration QA updates. * metric: register flow metrics only when we have a collector (#14529) the collector is absent when the pipeline is run in test with a NullMetricExt, or when the pipeline is explicitly configured to not collect metrics using `metric.collect: false`. * Unit tests and integration tests added for flow metrics. * Node stat spec and pipeline spec metric updates. * Metric keys statically imported, implicit error expectation added in metric spec. * Fix node status API spec after renaming flow metrics. * Removing flow metric from PipelinesInfo DS (used in peridoci metric snapshot), integration QA updates. * Rebasing with feature branch. * metric: register flow metrics only when we have a collector the collector is absent when the pipeline is run in test with a NullMetricExt, or when the pipeline is explicitly configured to not collect metrics using `metric.collect: false`. * Apply suggestions from code review Integration tests updated to test capturing the flow metrics. * Flow metrics expectation updated in tegration tests. * flow: refine integration expectations for reloads/monitoring Co-authored-by: Ry Biesemeyer <yaauie@users.noreply.github.com> Co-authored-by: Ry Biesemeyer <ry.biesemeyer@elastic.co> Co-authored-by: Mashhur <mashhur.sattorov@gmail.com>
* metric: add ScaledView with sub-unit precision to UptimeMetric By presenting a _view_ of our metric that maintains sub-unit precision, we prevent jitter that can be caused by our periodic poller not running at exactly our configured cadence. This is especially important as the UptimeMetric is used as the _denominator_ of several flow metrics, and a capture at 4.999s that truncates to 4s, causes the rate to be over-reported by ~25%. The `UptimeMetric.ScaledView` implements `Metric<Number>`, so its full lossless `BigDecimal` value is accessible to our `FlowMetric` at query time. * metrics: reduce window for too-frequent-captures bug and document it * fixup: provide mocked clock to flow metric
| pluginMetricsCounter = | ||
| LongCounter.fromRubyBase(pluginMetrics, MetricKeys.OUT_KEY); | ||
| pluginMetricsTime = LongCounter.fromRubyBase(pluginMetrics, PUSH_DURATION_KEY); | ||
| pluginMetricsTime = LongCounter.fromRubyBase(pluginMetrics, MetricKeys.PUSH_DURATION_KEY); |
There was a problem hiding this comment.
nit: we could static import MetricKeys .* to avoid repetitive usages, probability of conflict with other scopes is so low.
| this.nanoTimestamp = nanoTimestamp; | ||
| } | ||
|
|
||
| Double calculateRate(final Capture baseline) { |
There was a problem hiding this comment.
nit: we could use Optional since we are returning null values to avoid NPEs.
* flow metrics: code-style and readability pass * remove unused imports * cleanup: simplify usage of internal helpers * flow: migrate internals to use OptionalDouble
* flow: add global top-level flows * docs: add flow metrics
|
📃 DOCS PREVIEW ✨ https://logstash_14518.docs-preview.app.elstc.co/diff |
* Top level flow metrics unit tests added. * Add unit tests when config reloads, make sure top-level flow metrics didn't get reset. * Apply suggestions from code review Co-authored-by: Ry Biesemeyer <yaauie@users.noreply.github.com> * Validating against Hash test cases updated. * For the safety check against exact type in unit tests. Co-authored-by: Ry Biesemeyer <yaauie@users.noreply.github.com>
|
@mashhurs as communicated off-band, I'm passing this to you to do a final review and to merge once CI passes. We are intending to "freeze" what goes into this feature branch as of now, get it merged, and do any additional follow-up against |
|
📃 DOCS PREVIEW ✨ https://logstash_14518.docs-preview.app.elstc.co/diff |
karenzone
left a comment
There was a problem hiding this comment.
Good placement, good descriptions. Thank you for your work on this.
Release notes
Adds pipeline "flow" metrics to the node_stats API for each pipeline, which includes the
currentandlifetimerates for five key pipeline metrics:input_throughput,filter_throughput,output_throughput,queue_backpressure, andworker_concurrency.What does this PR do?
Implements Phase
0and1.Aof #14463, providing each pipeline'scurrentandlifetimerates for 5 key metrics:input_throughput:pipelines.*.events.in/ second of pipeline uptimefilter_throughput:pipelines.*.events.filtered/ second of pipeline uptimeoutput_throughput:pipelines.*.events.out/ second of pipeline uptimequeue_backpressure:pipelines.*.queue_push_duration/ second of pipeline uptimeworker_concurrency:duration/ second of pipeline uptimeThis Pull-Request is a combination of the efforts reviewed at a high-level in #14509 and #14514 for inclusion into this feature branch, and is a place for us to further review and document the now-integrated components.
Why is it important/What is the impact to the user?
As described in #14463, notably:
Checklist
Author's Checklist
uptime_in_millisto the API response.How to test this PR locally
pipelines.main.flowmetrics, and how they relate to the cumulative-value numbers they are tracking.