NEWS

Atlassian swaps metrics engine across 100,000 hosts without touching the code

Atlassian migrated its internal metrics platform to OpenTelemetry while keeping the StatsD interface intact for application teams

Atlassian swaps metrics engine across 100,000 hosts without touching the code
Image: Redação iMasters

Atlassian migrated its internal metrics platform to OpenTelemetry while keeping the StatsD interface intact for application teams. The case was published on October 1st. The platform receives data from around 100,000 hosts across 14 regions, under a 99.95% SLO.

The strategy deserves attention. The company rebuilt collection, routing, aggregation, and forwarding underneath, while the contract with applications stayed the same.

Therefore, no team needed to reinstrument services before the switch.

Why gostatsd stopped being enough

The company had maintained gostatsd, an open-source implementation of StatsD, for most of the last decade. This component handled collection as an auxiliary process on each host, as well as aggregation.

The limit showed up in the architecture. Based on UDP, it didn't support traces and logs.

Meanwhile, much of the internal telemetry was already migrating to OpenTelemetry. This project, incidentally, graduated as a CNCF project in May this year, with metrics, traces, and logs as its core signals.

Atlassian split the pipeline into four stages

The new design separates collection, ingestion, aggregation, and forwarding. Each stage uses a specific distribution of the OpenTelemetry Collector.

This way, teams can replace components independently.

The collection layer accepts data in both StatsD and OTLP at the same time. Consequently, older applications keep sending over UDP, while new instrumentation sends directly in OpenTelemetry.

The unification brought an immediate gain. By combining metrics and tracing in the same Collector-based sidecar, the company removed the separate gostatsd sidecar.

The result reached about 3.9% less CPU per service on the most expensive Micros services. This represented approximately a 30% reduction in sidecar cost across the fleet.

The routing problem that stream ID solved

Here's the most interesting part of the case. Metric aggregation is a stateful operation, so points from the same time series need to reach the same aggregator.

The previous solution used an internal proxy called nomad. It distributed traffic to shards by hashing service and environment.

The design created imbalance. Since volume varies widely between services, a shard responsible for a large service received much more load.

The switch came through the Collector's load balancing exporter. Instead of routing by service, it uses a stream ID that identifies an individual time series.

This way, different series from the same service spread across multiple shards. At the same time, each series stays on a consistent shard.

The company reported a more uniform CPU distribution across shards. In addition, the aggregation pool started scaling down more during periods of low activity.

The aggregation numbers are impressive

The volume gives a sense of the challenge's scale. The pipeline receives around 4.8 billion points per minute.

After aggregation, storage keeps around 220 million. In other words, the reduction reaches approximately 96%.

There was also a specific technical obstacle. Much of the data uses delta temporality, where each point represents the change since the previous measurement.

Existing upstream components aggregated these deltas differently than needed. For that reason, the team developed its own delta aggregation processor and published it through Atlassian Labs.

The new aggregation layer consumes about half the CPU of the previous system under the same traffic.

Forwarding, Lambda, and the gradual rollout

The forwarding stage became a stateless distribution called metrics gateway. It sends telemetry to destinations such as SignalFx and Amazon S3, using native retry, queueing, and backpressure features.

For workloads on AWS Lambda, where a conventional sidecar can't run, the company created its own extension. It keeps the same StatsD address and the same environment variables as the previous implementation.

The rollout followed a staircase pattern. After tests in development, staging, and less critical workloads, production rose from about 1% to 10%, 50%, and full coverage.

Both systems ran side by side for most of the migration. So, operational parity became a requirement along the way.

It's worth noting a lesson about measurement. According to the team, benchmarks alone failed to reveal the real cost, and continuous profiling in production pointed to where to optimize.

What comes next

The old components still weigh heavily. gostatsd aggregators and the nomad service account for about 38% of CPU requests in the metrics clusters, with nomad alone responsible for approximately 13%.

The next step moves instrumentation to the application side. The goal is to move away from StatsD, DogStatsD, and internal clients toward OpenTelemetry SDKs.

This stage comes after the infrastructure switch. This way, each service migrates at its own pace.

Follow our profile on Instagram!

Translated from the Brazilian Portuguese original · Read the original

More from Redação iMasters
View profile →