Repository navigation
fix(aggregators): avoid holding RunningAggregator lock during metric delivery in Push() - #19786
spaghettidba wants to merge 1 commit into
Conversation
…delivery in Push() RunningAggregator.Push() previously called the aggregator plugin's Push() directly against the real telegraf.Accumulator while holding the same mutex used by Add(). The real accumulator's AddMetric/AddFields calls block on a bounded channel to the next processing stage; if that channel is momentarily full (e.g. a slow output, or a large/ high-cardinality aggregation taking non-trivial time to iterate and deliver), the lock could be held long enough to starve concurrent Add() calls. Because input plugins such as influxdb_listener call Add() (via the accumulator) synchronously while handling an inbound write request, prolonged lock contention here can cause those requests to stall or time out client-side, resulting in original metrics being lost before they ever reach Telegraf -- even with drop_original = false, since the original metric never made it into the pipeline in the first place. Fix Push() to write into an in-memory collector while holding the lock, then release the lock before delivering the buffered metrics to the real accumulator. MakeMetric/precision handling still happens exactly once, inside acc.AddMetric, unchanged from before. Added regression tests demonstrating that Add() is no longer blocked by a slow/blocking Push() and that no metrics are lost during delivery.
|
Download PR build artifacts for linux_amd64.tar.gz, darwin_arm64.tar.gz, and windows_amd64.zip. 📦 Click here to get additional PR build artifactsArtifact URLs |
|
@spaghettidba please restore the PR description template, especially the AI section, as we cannot review your contribution otherwise. |
|
Sorry about that, copilot opened the pr for me and apparently didn't respect the template. |
|
No worries @spaghettidba! Does this PR fix an open issue? If so, please add this part to the PR description too so the issue is automatically closed when we merge this PR. |
|
No, it does not close an open issue. I was trying to fix a different problem (points outside of window being dropped by aggregator despite |
|
@spaghettidba Thanks for digging into this. Before we change the aggregator core I'd like to understand where your stall actually comes from. Inputs don't call RunningAggregator.Add() themselves, a single agent goroutine does, so what the listener sees is back-pressure from that goroutine waiting on the lock. The PR only moves the hand-off to the next stage out of the lock, the plugin's own Push() still runs under it. I measured a starlark aggregator with 40k groups and a plain output downstream. The plugin's Push() took about 26ms and handing the metrics on took about 8ms, so the part this PR moves is the small one. To get to 300ms stalls something downstream must be slow, like processors running after the aggregator or a disk buffer. Could you open an issue with your config and the internal_aggregate push_time_ns values you see? That stat currently includes the hand-off, so it shows directly whether delivery is your bottleneck. If it is, I'm fine with the approach, but let's keep the second accumulator as small as possible since it has to stay in line with agent/accumulator.go. |
Summary
RunningAggregator.Push()calls the aggregator plugin'sPush()method directly against the realtelegraf.Accumulator, while holding the samesync.MutexthatAdd()uses. The real accumulator'sAddMetric/AddFieldscalls ultimately block on a bounded channel to the next processing stage (seeagent/accumulator.go). If that channel is momentarily full — e.g. a slow output, or an aggregation with high tag/group cardinality that takes non-trivial time to iterate and deliver — the lock can be held long enough to starve concurrentAdd()calls.Because input plugins such as
influxdb_listenercallAdd()synchronously while handling an inbound HTTP write request, prolonged lock contention here can cause those requests to stall or time out client-side. The result is original metrics being lost before they ever reach Telegraf, even withdrop_original = false, since the metric never made it into the pipeline to begin with. This is silent from Telegraf's point of view (no error is logged), which made it especially hard to diagnose in production.Reproduction
I reproduced this with a live A/B test: two Telegraf processes with an identical
influxdb_listenerinput andoutputs.influxdboutput, one with a high-cardinalityaggregators.starlarkaggregator attached and one without (control), both under a background feeder generating ~40,000 distinct tag combinations. A client issuing HTTP writes with a strict 300ms timeout saw periodic timeouts against the aggregator-enabled instance (correlated withperiodrollovers), and zero timeouts against the control instance. After the fix, timeouts against the aggregator-enabled instance dropped to zero as well.Fix
Push()now writes into an in-memorypushCollector(implementstelegraf.Accumulator, just appends to a slice) while holding the lock. Once the plugin'sPush()/Reset()calls are done, the lock is released, and then the buffered metrics are delivered to the real accumulator viaacc.AddMetric().MakeMetric/precision handling still happens exactly once, insideacc.AddMetric, unchanged from before. ApushTrackingCollectorvariant supports the (rare) case of a plugin callingacc.WithTracking()from withinPush().Testing
TestRunningAggregatorAddNotBlockedByPushandTestRunningAggregatorPushDeliversAllMetricsWithoutHoldingLocktomodels/running_aggregator_test.go, using aslowPushAggregator(blocks insidePush) and ablockingAccumulator(gatesAddFields/AddMetricon a channel) to simulate a slow downstream accumulator.go test ./models/...andgo test ./agent/...pass.go vet ./models/...andgofmt -lare clean.-racecould not be run in my environment due to lack of a C compiler/CGO; would appreciate CI running it.)Checklist