Skip to content

Feature: flow metrics integration - #14518

Merged
mashhurs merged 12 commits into
mainfrom
feature/flow-metrics-integration
Sep 19, 2022
Merged

mashhurs merged 12 commits into
mainfrom
feature/flow-metrics-integration

Conversation

@yaauie

@yaauie yaauie commented Sep 9, 2022 •

Copy link
Copy Markdown
Member

Release notes

Adds pipeline "flow" metrics to the node_stats API for each pipeline, which includes the current and lifetime rates for five key pipeline metrics: input_throughput, filter_throughput, output_throughput, queue_backpressure, and worker_concurrency.

What does this PR do?

Implements Phase 0 and 1.A of #14463, providing each pipeline's current and lifetime rates for 5 key metrics:

  • input_throughput: pipelines.*.events.in / second of pipeline uptime
  • filter_throughput: pipelines.*.events.filtered / second of pipeline uptime
  • output_throughput: pipelines.*.events.out / second of pipeline uptime
  • queue_backpressure: pipelines.*.queue_push_duration / second of pipeline uptime
  • worker_concurrency: duration / second of pipeline uptime

This 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:

Problem Statement

It is often difficult to understand the health of a pipeline, including whether it is exerting or propagating back-pressure or otherwise staying reasonably “caught up” with its inputs. The cumulative-value metrics exposed by the API are not typically useful on their own, and when we are debugging a pipeline we often require a pair of captures and manual math in order to see the flow of events through the pipeline.

We want to provide better visibility into the flow of events through a Logstash pipeline in order to easier pinpoint the source of degraded status.

Scope and Goals

We will provide visibility into the current, recent, and lifetime flow rates of pipeline-level metrics in the /_node/stats API, so that our users can have more contextual information when analyzing the throughput of a pipeline.

Checklist

  • My code follows the style guidelines of this project
  • I have commented my code, particularly in hard-to-understand areas
  • I have made corresponding changes to the documentation
  • I have made corresponding change to the default configuration files (and/or docker env variables)
  • I have added tests that prove my fix is effective or that my feature works

Author's Checklist

  • periodic poller warns if it is unable to keep up with configured polling frequency
  • expose our pipeline's new uptime_in_millis to the API response.

How to test this PR locally

  1. Run a pipeline that handles jittery input data.
    ruby -e '100.times.map { |j| Thread.new(j) { |k| 100000.times { |i| $stdout.write("ok[#{k}/#{i}]\n"); (i%17).zero? && sleep(Random.rand(10)) } } }.map(&:join)' | bin/logstash -e 'input { stdin {} } output { sink {} }'
    
  2. once the pipeline is started, query the node stats API, and observe the new pipelines.main.flow metrics, and how they relate to the cumulative-value numbers they are tracking.

yaauie and others added 2 commits September 9, 2022 14:36
* 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 yaauie left a comment

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

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

First pass: concurrency could be made more clear as worker_concurrency, and backpressure similarly reads more clearly as queue_backpressure.

Comment thread logstash-core/lib/logstash/api/commands/stats.rb Outdated
Comment thread logstash-core/spec/logstash/api/commands/stats_spec.rb Outdated
Comment thread logstash-core/spec/logstash/api/modules/node_stats_spec.rb Outdated
Comment thread logstash-core/spec/logstash/instrument/wrapped_write_client_spec.rb Outdated
Comment thread logstash-core/src/main/java/org/logstash/execution/AbstractPipelineExt.java Outdated
Comment thread tools/benchmark-cli/src/test/resources/org/logstash/benchmark/cli/metrics.json Outdated
Comment thread logstash-core/lib/logstash/api/commands/stats.rb Outdated
Comment thread logstash-core/src/main/java/org/logstash/execution/AbstractPipelineExt.java Outdated
mashhurs and others added 5 commits September 13, 2022 07:01
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
@yaauie
yaauie marked this pull request as ready for review September 15, 2022 22:04
pluginMetricsCounter =
LongCounter.fromRubyBase(pluginMetrics, MetricKeys.OUT_KEY);
pluginMetricsTime = LongCounter.fromRubyBase(pluginMetrics, PUSH_DURATION_KEY);
pluginMetricsTime = LongCounter.fromRubyBase(pluginMetrics, MetricKeys.PUSH_DURATION_KEY);

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: 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) {

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: we could use Optional since we are returning null values to avoid NPEs.

@yaauie yaauie mentioned this pull request Sep 16, 2022
2 tasks done
* 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
@yaauie
yaauie requested a review from karenzone September 19, 2022 16:58
@github-actions

Copy link
Copy Markdown
Contributor

📃 DOCS PREVIEW ✨ https://logstash_14518.docs-preview.app.elstc.co/diff

mashhurs and others added 2 commits September 19, 2022 13:20
* 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>
@yaauie
yaauie requested a review from mashhurs September 19, 2022 20:25
@yaauie

yaauie commented Sep 19, 2022

Copy link
Copy Markdown
Member Author

@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 main directly.

@github-actions

Copy link
Copy Markdown
Contributor

📃 DOCS PREVIEW ✨ https://logstash_14518.docs-preview.app.elstc.co/diff

@karenzone karenzone 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.

Good placement, good descriptions. Thank you for your work on this.

@mashhurs mashhurs 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.

💯 We have worked on this feature branch and every small changes were reviewed and then merged into feature branch. Once more went trough all changes and all seem reasonable, as expected. Thanks @yaauie making this happen!

@mashhurs
mashhurs merged commit 6e0b365 into main Sep 19, 2022
@jsvd
jsvd deleted the feature/flow-metrics-integration branch September 20, 2022 13:35
@jsvd
jsvd restored the feature/flow-metrics-integration branch September 20, 2022 13:35
@jsvd jsvd added the v8.5.0 label Oct 6, 2022
@jsvd
jsvd deleted the feature/flow-metrics-integration branch June 2, 2023 12:11
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

5 participants