Skip to content

fix(aggregators): avoid holding RunningAggregator lock during metric delivery in Push() - #19786

Open
spaghettidba wants to merge 1 commit into
influxdata:masterfrom
spaghettidba:fix/aggregator-push-lock-contention
Open

spaghettidba wants to merge 1 commit into
influxdata:masterfrom
spaghettidba:fix/aggregator-push-lock-contention

Conversation

@spaghettidba

@spaghettidba spaghettidba commented Sep 25, 2026 •

Copy link
Copy Markdown
Contributor

Summary

RunningAggregator.Push() calls the aggregator plugin's Push() method directly against the real telegraf.Accumulator, while holding the same sync.Mutex that Add() uses. The real accumulator's AddMetric/AddFields calls ultimately block on a bounded channel to the next processing stage (see agent/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 concurrent Add() calls.

Because input plugins such as influxdb_listener call Add() 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 with drop_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_listener input and outputs.influxdb output, one with a high-cardinality aggregators.starlark aggregator 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 with period rollovers), 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-memory pushCollector (implements telegraf.Accumulator, just appends to a slice) while holding the lock. Once the plugin's Push()/Reset() calls are done, the lock is released, and then the buffered metrics are delivered to the real accumulator via acc.AddMetric(). MakeMetric/precision handling still happens exactly once, inside acc.AddMetric, unchanged from before. A pushTrackingCollector variant supports the (rare) case of a plugin calling acc.WithTracking() from within Push().

Testing

  • Added TestRunningAggregatorAddNotBlockedByPush and TestRunningAggregatorPushDeliversAllMetricsWithoutHoldingLock to models/running_aggregator_test.go, using a slowPushAggregator (blocks inside Push) and a blockingAccumulator (gates AddFields/AddMetric on a channel) to simulate a slow downstream accumulator.
  • Verified the new tests fail against the pre-fix code and pass against the fix.
  • go test ./models/... and go test ./agent/... pass.
  • go vet ./models/... and gofmt -l are clean.
  • (Note: -race could not be run in my environment due to lack of a C compiler/CGO; would appreciate CI running it.)

Checklist

…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.
@telegraf-tiger telegraf-tiger Bot added the fix pr to fix corresponding bug label Sep 25, 2026
@telegraf-tiger

Copy link
Copy Markdown
Contributor

Download PR build artifacts for linux_amd64.tar.gz, darwin_arm64.tar.gz, and windows_amd64.zip.
Downloads for additional architectures and packages are available below.

⚠️ This pull request increases the Telegraf binary size by 5.66 % for linux amd64 (new size: 325.7 MB, nightly size 308.2 MB)

📦 Click here to get additional PR build artifacts

Artifact URLs

. DEB . RPM . TAR . GZ . ZIP
amd64.deb aarch64.rpm darwin_amd64.tar.gz windows_amd64.zip
arm64.deb armel.rpm darwin_arm64.tar.gz windows_arm64.zip
armel.deb armv6hl.rpm freebsd_amd64.tar.gz windows_i386.zip
armhf.deb i386.rpm freebsd_armv7.tar.gz
i386.deb ppc64le.rpm freebsd_i386.tar.gz
mips.deb riscv64.rpm linux_amd64.tar.gz
mipsel.deb s390x.rpm linux_arm64.tar.gz
ppc64el.deb x86_64.rpm linux_armel.tar.gz
riscv64.deb linux_armhf.tar.gz
s390x.deb linux_i386.tar.gz
linux_mips.tar.gz
linux_mipsel.tar.gz
linux_ppc64le.tar.gz
linux_riscv64.tar.gz
linux_s390x.tar.gz

@srebhan

srebhan commented Sep 28, 2026

Copy link
Copy Markdown
Member

@spaghettidba please restore the PR description template, especially the AI section, as we cannot review your contribution otherwise.

@srebhan srebhan self-assigned this Sep 28, 2026
@srebhan srebhan added plugin/aggregator 1. Request for new aggregator plugins 2. Issues/PRs that are related to aggregator plugins waiting for response waiting for response from contributor labels Sep 28, 2026
@spaghettidba

Copy link
Copy Markdown
Contributor Author

Sorry about that, copilot opened the pr for me and apparently didn't respect the template.
I think I restored it, please tell me if something is still missing.
Regarding the pr, I think it is a very unlikely edge case, but I thought it was worth reporting

@telegraf-tiger telegraf-tiger Bot removed the waiting for response waiting for response from contributor label Sep 28, 2026
@srebhan

srebhan commented Sep 30, 2026

Copy link
Copy Markdown
Member

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.

@spaghettidba

Copy link
Copy Markdown
Contributor Author

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 drop_original = false) and came across this potential problem

@skartikey

Copy link
Copy Markdown
Contributor

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

@srebhan srebhan assigned skartikey and unassigned srebhan Oct 7, 2026
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

fix pr to fix corresponding bug plugin/aggregator 1. Request for new aggregator plugins 2. Issues/PRs that are related to aggregator plugins

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants