Bound Prometheus distribution buffering
This commit is contained in:
parent
695f6ebc0f
commit
882df25aa9
|
|
@ -847,6 +847,26 @@ integer-based; divide by `1_000_000` in PromQL when seconds are required.
|
||||||
Definitions intentionally have no request path, user, request, or event-name
|
Definitions intentionally have no request path, user, request, or event-name
|
||||||
labels that could create unbounded cardinality.
|
labels that could create unbounded cardinality.
|
||||||
|
|
||||||
|
The locked `telemetry_metrics_prometheus_core` reporter aggregates
|
||||||
|
distribution samples only when a scrape occurs. Each application VM therefore
|
||||||
|
also runs a supervised internal aggregation every ten seconds. This keeps the
|
||||||
|
reporter's raw distribution table independent of whether an external
|
||||||
|
Prometheus server is currently configured, while preserving the cumulative
|
||||||
|
histograms exposed by the authenticated endpoint. The generated exposition
|
||||||
|
text from the internal aggregation is discarded.
|
||||||
|
|
||||||
|
After a deployment, the raw table can be inspected without exposing metric
|
||||||
|
credentials:
|
||||||
|
|
||||||
|
```bash
|
||||||
|
/app/bin/who_need_help eval \
|
||||||
|
'IO.inspect(:ets.info(:prometheus_metrics_dist, :size), label: "pending_distribution_samples")'
|
||||||
|
```
|
||||||
|
|
||||||
|
Under idle conditions the value returns to zero after the next aggregation
|
||||||
|
interval. During traffic it represents only samples received since the most
|
||||||
|
recent internal or external scrape; it is not the cumulative histogram count.
|
||||||
|
|
||||||
Metrics are local to each BEAM process. Discover and scrape every web pod or
|
Metrics are local to each BEAM process. Discover and scrape every web pod or
|
||||||
container as a distinct target and preserve Prometheus's `instance` label. A
|
container as a distinct target and preserve Prometheus's `instance` label. A
|
||||||
request through the load-balanced public route reaches only one replica and is
|
request through the load-balanced public route reaches only one replica and is
|
||||||
|
|
|
||||||
35
lib/who_need_help_web/prometheus_distribution_drain.ex
Normal file
35
lib/who_need_help_web/prometheus_distribution_drain.ex
Normal file
|
|
@ -0,0 +1,35 @@
|
||||||
|
defmodule WhoNeedHelpWeb.PrometheusDistributionDrain do
|
||||||
|
@moduledoc false
|
||||||
|
|
||||||
|
use GenServer
|
||||||
|
|
||||||
|
@default_interval 10_000
|
||||||
|
|
||||||
|
def start_link(opts) do
|
||||||
|
reporter_name = Keyword.fetch!(opts, :reporter_name)
|
||||||
|
interval = Keyword.get(opts, :interval, @default_interval)
|
||||||
|
|
||||||
|
GenServer.start_link(__MODULE__, {reporter_name, interval})
|
||||||
|
end
|
||||||
|
|
||||||
|
@impl true
|
||||||
|
def init({reporter_name, interval}) when is_integer(interval) and interval > 0 do
|
||||||
|
schedule_drain(interval)
|
||||||
|
{:ok, %{reporter_name: reporter_name, interval: interval}}
|
||||||
|
end
|
||||||
|
|
||||||
|
@impl true
|
||||||
|
def handle_info(:drain, state) do
|
||||||
|
# telemetry_metrics_prometheus_core buffers every distribution sample until
|
||||||
|
# scrape time. Aggregate on a bounded interval even when no external
|
||||||
|
# Prometheus server is configured.
|
||||||
|
_scrape = TelemetryMetricsPrometheus.Core.scrape(state.reporter_name)
|
||||||
|
|
||||||
|
schedule_drain(state.interval)
|
||||||
|
{:noreply, state}
|
||||||
|
end
|
||||||
|
|
||||||
|
defp schedule_drain(interval) do
|
||||||
|
Process.send_after(self(), :drain, interval)
|
||||||
|
end
|
||||||
|
end
|
||||||
|
|
@ -12,7 +12,9 @@ defmodule WhoNeedHelpWeb.Telemetry do
|
||||||
if Application.fetch_env!(:who_need_help, :app_role) in [:web, :worker, :combined] do
|
if Application.fetch_env!(:who_need_help, :app_role) in [:web, :worker, :combined] do
|
||||||
[
|
[
|
||||||
{TelemetryMetricsPrometheus.Core,
|
{TelemetryMetricsPrometheus.Core,
|
||||||
name: :prometheus_metrics, metrics: prometheus_metrics(), start_async: false}
|
name: :prometheus_metrics, metrics: prometheus_metrics(), start_async: false},
|
||||||
|
{WhoNeedHelpWeb.PrometheusDistributionDrain,
|
||||||
|
reporter_name: :prometheus_metrics, interval: 10_000}
|
||||||
]
|
]
|
||||||
else
|
else
|
||||||
[]
|
[]
|
||||||
|
|
|
||||||
|
|
@ -0,0 +1,57 @@
|
||||||
|
defmodule WhoNeedHelpWeb.PrometheusDistributionDrainTest do
|
||||||
|
use ExUnit.Case, async: false
|
||||||
|
|
||||||
|
import Telemetry.Metrics
|
||||||
|
|
||||||
|
alias TelemetryMetricsPrometheus.Core.Registry
|
||||||
|
alias WhoNeedHelpWeb.PrometheusDistributionDrain
|
||||||
|
|
||||||
|
@reporter_name :prometheus_distribution_drain_test
|
||||||
|
@event_name [:who_need_help, :test, :distribution_drain]
|
||||||
|
|
||||||
|
test "aggregates distribution samples without an external metrics scrape" do
|
||||||
|
metric =
|
||||||
|
distribution("who_need_help.test.distribution.seconds",
|
||||||
|
event_name: @event_name,
|
||||||
|
measurement: :duration,
|
||||||
|
reporter_options: [buckets: [0.1, 0.5, 1.0]]
|
||||||
|
)
|
||||||
|
|
||||||
|
start_supervised!(
|
||||||
|
{TelemetryMetricsPrometheus.Core,
|
||||||
|
name: @reporter_name, metrics: [metric], start_async: false}
|
||||||
|
)
|
||||||
|
|
||||||
|
config = Registry.config(@reporter_name)
|
||||||
|
|
||||||
|
for duration <- 1..2_000 do
|
||||||
|
:telemetry.execute(@event_name, %{duration: duration / 10_000}, %{})
|
||||||
|
end
|
||||||
|
|
||||||
|
assert :ets.info(config.dist_table_id, :size) == 2_000
|
||||||
|
|
||||||
|
start_supervised!({PrometheusDistributionDrain, reporter_name: @reporter_name, interval: 10})
|
||||||
|
|
||||||
|
assert_eventually(fn -> :ets.info(config.dist_table_id, :size) == 0 end)
|
||||||
|
|
||||||
|
body = TelemetryMetricsPrometheus.Core.scrape(@reporter_name)
|
||||||
|
|
||||||
|
assert body =~ "who_need_help_test_distribution_seconds_count 2000"
|
||||||
|
assert :ets.info(config.dist_table_id, :size) == 0
|
||||||
|
end
|
||||||
|
|
||||||
|
defp assert_eventually(fun, attempts \\ 100)
|
||||||
|
|
||||||
|
defp assert_eventually(fun, attempts) when attempts > 0 do
|
||||||
|
if fun.() do
|
||||||
|
assert true
|
||||||
|
else
|
||||||
|
Process.sleep(10)
|
||||||
|
assert_eventually(fun, attempts - 1)
|
||||||
|
end
|
||||||
|
end
|
||||||
|
|
||||||
|
defp assert_eventually(_fun, 0) do
|
||||||
|
flunk("condition did not become true before the timeout")
|
||||||
|
end
|
||||||
|
end
|
||||||
Loading…
Reference in New Issue
Block a user