Skip to content

Build1 publisher3 min readPublished

Atlassian rebuilds its 100,000-host metrics pipeline on OpenTelemetry behind the old StatsD interface

Atlassian rebuilt a metrics platform fed by about 100,000 hosts on OpenTelemetry while its applications kept sending StatsD packets to the same address. The platform team did the work alone. Moving thousands of services onto OpenTelemetry SDKs comes later.

The Engineer · Build desk

Illustration accompanying Atlassian rebuilds its 100,000-host metrics pipeline on OpenTelemetry behind the old StatsD interface

What happened

  • Atlassian left gostatsd because its UDP-only design could not carry traces or logs, and staying would have meant rebuilding features the OpenTelemetry Collector community was already developing.
  • The rebuilt platform runs four purpose-built Collector distributions, one each for collection, ingest, aggregation and forwarding, so one stage can change without replacing the rest.
  • The platform receives about 4.8 billion data points a minute and stores about 220 million of them after aggregation.
  • A stateless Collector distribution called metrics-gateway replaced a bespoke forwarding service, fanning data out to destinations including SignalFx and S3.
  • Rollout started with development and staging workloads and chosen early adopters, then widened through 1%, 10%, 50% and 100% stages.

Compiled by The EngineerSomething wrong?How this is made

Why it matters

  • decision Platform teams on StatsD now have a tested order of operations: replace the backend behind the existing wire format first, then run the SDK move as a separate project on its own schedule.
  • capability Teams whose StatsD traffic produces delta metrics can evaluate Atlassian's open-sourced aggregation processor before writing their own.
  • capability With the Collector handling retry, queueing and backpressure, sending metrics to one more backend is a configuration change, so dual-writing during a vendor switch is cheaper to try.

Applications at Atlassian send StatsD packets over UDP to a fixed address, and they still do [7]. The software listening on that address changed. Collection now runs in the OpenTelemetry Collector distribution the tracing team already deployed as a sidecar. An OTLP receiver sits next to the StatsD one for newer services that send OpenTelemetry metrics [7]. Serverless workloads cannot run a sidecar, so an OpenTelemetry Lambda extension presents the same interface there [9]. The authors say the team avoided asking thousands of services to switch from StatsD clients to the OpenTelemetry SDK straight away [4].

The ingest tier holds the constraint that shapes everything after it. Metric aggregation is stateful, so every point for a given time series must reach the same aggregator [10]. The old nomad proxy met that rule by hashing each service and environment to one shard. The biggest services piled onto a few busy replicas [10]. The replacement uses the OpenTelemetry contrib load-balancing exporter and hashes by stream ID [11]. In my view stream ID is the right key. A single series is the smallest unit an aggregator has to see whole, and hashing on it lets one large service spread across the pool while each of its series stays together [11]. Atlassian reports more even CPU, fewer hot shards and better off-peak scaling [11].

At the reported inflow, the platform takes in about 80 million data points a second [1]. Spread evenly, that would be about 48,000 points a minute from each of the 100,000 hosts across 14 regions [2]. The rounded published figures give a reduction of about 95.4%, against the roughly 96% stated [3]. Upstream components did not aggregate delta metrics the way Atlassian needed, so it built its own aggregation processor and open-sourced it [13]. The revised aggregation tier uses about half the CPU for the same traffic [13]. "Small tests and benchmarks were not enough; continuous profiling in prod is what actually told us where to optimise," the authors wrote [14].

The sidecar savings are figures from Atlassian's own fleet. Folding metrics into the tracing sidecar saved an average of 3.9% CPU for each of its most expensive Micros services, and the team estimates sidecar cost fell about 30% across the fleet [8]. The saving comes from retiring the gostatsd sidecar and running one collector where two processes ran before [7][8]. It transfers to a fleet that already runs a Collector sidecar for traces beside a separate StatsD agent. A fleet with no tracing sidecar has nothing to fold metrics into, and would be adding a sidecar to save on sidecars.

"We kept the interface and rebuilt everything behind it, which turned an org-wide migration into a platform-team migration," Iris Grace Endozo, Farzad Vazirnia and Albert Kerr wrote on the CNCF blog [5]. That holds for the pipeline. The application side has been pushed to a later phase. Atlassian's next planned step moves instrumentation from StatsD, DogStatsD and vendor clients to OpenTelemetry SDKs, the change the kept interface let thousands of services put off [17][4]. The pipeline is not finished either. Gostatsd aggregation and nomad still account for about 38% of CPU requests in the metrics clusters, and nomad alone uses about 13% of total resources [16]. The end-to-end OpenTelemetry pipeline is complete only once both are removed [16]. The published account does not report how the 99.95% SLO held during the staged rollouts [2].

What to watch

  • Figures after gostatsd aggregation and nomad are removed, showing how much of their roughly 38% share of metrics-cluster CPU requests comes back as savings.
  • Whether the move of application instrumentation to OpenTelemetry SDKs stays with the platform team or turns into a per-service project for thousands of owners.
  • Whether Atlassian's open-sourced delta aggregation processor moves upstream into OpenTelemetry contrib, taking a custom component off its own maintenance list.
Loading claim ledger
Loading source directory links
Loading share composer
Loading topic controls
Loading related stories