Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 9 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -71,6 +71,15 @@ Increment the:
at shutdown
[#4637](https://github.com/open-telemetry/opentelemetry-cpp/pull/4637)

* [SDK] Fix spatial re-aggregation for asynchronous instruments. A view which
drops attributes is now applied to asynchronous instruments as well, and the
observations which collapse onto the same attribute set are re-aggregated
(summed up for the additive instruments) instead of the last observation
overwriting the previous ones.
Note that `AsyncMetricStorage`'s constructor now takes the view's
`AttributesProcessor`, matching `SyncMetricStorage`.
[#1724](https://github.com/open-telemetry/opentelemetry-cpp/issues/1724)

## [1.29.0] 2026-09-13

* [RELEASE] Bump main branch to 1.29.0-dev (#4259)
Expand Down
179 changes: 158 additions & 21 deletions sdk/include/opentelemetry/sdk/metrics/state/async_metric_storage.h
Original file line number Diff line number Diff line change
Expand Up @@ -3,12 +3,15 @@

#pragma once

#include <functional>
#include <memory>
#include <mutex>
#include <unordered_map>
#include <utility>

#include "opentelemetry/nostd/shared_ptr.h"
#include "opentelemetry/sdk/common/attributemap_hash.h"
#include "opentelemetry/sdk/metrics/aggregation/aggregation.h"
#include "opentelemetry/sdk/metrics/aggregation/aggregation_config.h"
#include "opentelemetry/sdk/metrics/aggregation/default_aggregation.h"

Expand Down Expand Up @@ -37,6 +40,7 @@ class AsyncMetricStorage : public MetricStorage, public AsyncWritableMetricStora
public:
AsyncMetricStorage(const InstrumentDescriptor &instrument_descriptor,
const AggregationType aggregation_type,
std::shared_ptr<const AttributesProcessor> attributes_processor,
#ifdef ENABLE_METRICS_EXEMPLAR_PREVIEW
ExemplarFilterType exemplar_filter_type,
nostd::shared_ptr<ExemplarReservoir> &&exemplar_reservoir,
Expand All @@ -45,24 +49,34 @@ class AsyncMetricStorage : public MetricStorage, public AsyncWritableMetricStora
: instrument_descriptor_(instrument_descriptor),
aggregation_type_{aggregation_type},
aggregation_config_{AggregationConfig::GetOrDefault(aggregation_config)},
cumulative_hash_map_(
attributes_processor_{std::move(attributes_processor)},
is_monotonic_sum_{IsMonotonicSum(aggregation_type, instrument_descriptor)},
last_observed_hash_map_(
std::make_unique<AttributesHashMap>(aggregation_config_->cardinality_limit_)),
delta_hash_map_(
std::make_unique<AttributesHashMap>(aggregation_config_->cardinality_limit_)),
round_hash_map_(
std::make_unique<AttributesHashMap>(aggregation_config_->cardinality_limit_)),
#ifdef ENABLE_METRICS_EXEMPLAR_PREVIEW
exemplar_filter_type_(exemplar_filter_type),
exemplar_reservoir_(std::move(exemplar_reservoir)),
#endif
temporal_metric_storage_(instrument_descriptor, aggregation_type, aggregation_config)
{}
{
create_default_aggregation_ = [this]() -> std::unique_ptr<Aggregation> {
return DefaultAggregation::CreateAggregation(aggregation_type_, instrument_descriptor_);
};
}

/**
* Converts the absolute values reported by the callbacks into the deltas the temporal storage
* consumes. Runs once per callback registered on the instrument, so what it accumulates stays
* until the next Collect().
*/
template <class T>
void Record(const std::unordered_map<MetricAttributes, T, AttributeHashGenerator> &measurements,
opentelemetry::common::SystemTimestamp /* observation_time */) noexcept
{
// Async counter always record monotonically increasing values, and the
// exporter/reader can request either for delta or cumulative value.
// So we convert the async counter value to delta before passing it to temporal storage.
std::lock_guard<std::mutex> guard(hashmap_lock_);
#ifdef ENABLE_METRICS_EXEMPLAR_PREVIEW
const bool offer_exemplars =
Expand All @@ -76,25 +90,15 @@ class AsyncMetricStorage : public MetricStorage, public AsyncWritableMetricStora
exemplar_reservoir_->OfferMeasurement(measurement.second, measurement.first, {});
}
#endif

auto aggr = DefaultAggregation::CreateAggregation(aggregation_type_, instrument_descriptor_);
aggr->Aggregate(measurement.second);
auto prev = cumulative_hash_map_->Get(measurement.first);
if (prev)
if (is_monotonic_sum_)
{
auto delta = prev->Diff(*aggr);
// store received value in cumulative map, and the diff in delta map (to pass it to temporal
// storage)
cumulative_hash_map_->Set(measurement.first, std::move(aggr));
delta_hash_map_->Set(measurement.first, std::move(delta));
RecordMonotonicSum(measurement.first, measurement.second);
}
else
{
// store received value in cumulative and delta map.
cumulative_hash_map_->Set(
measurement.first,
DefaultAggregation::CloneAggregation(aggregation_type_, instrument_descriptor_, *aggr));
delta_hash_map_->Set(measurement.first, std::move(aggr));
round_hash_map_
->GetOrSetDefault(FilterAttributes(measurement.first), create_default_aggregation_)
->Aggregate(measurement.second);
}
}
}
Expand Down Expand Up @@ -131,6 +135,10 @@ class AsyncMetricStorage : public MetricStorage, public AsyncWritableMetricStora
std::shared_ptr<AttributesHashMap> delta_metrics = nullptr;
{
std::lock_guard<std::mutex> guard(hashmap_lock_);
if (!is_monotonic_sum_)
{
BuildDeltaFromRound();
}
delta_metrics = std::move(delta_hash_map_);
delta_hash_map_ =
std::make_unique<AttributesHashMap>(aggregation_config_->cardinality_limit_);
Expand All @@ -143,11 +151,140 @@ class AsyncMetricStorage : public MetricStorage, public AsyncWritableMetricStora
}

private:
/**
* Differences a monotonic sum per source series - keyed by the unfiltered attributes - before
* the view merges anything. Diffing the merged group instead would turn a series which stops
* being reported into a spurious decrease.
*/
template <class T>
void RecordMonotonicSum(const MetricAttributes &attributes, T value)
{
auto observed = create_default_aggregation_();
observed->Aggregate(value);

std::unique_ptr<Aggregation> delta;
auto previous = last_observed_hash_map_->Get(attributes);
if (previous)
{
delta = previous->Diff(*observed);
}
else
{
delta = DefaultAggregation::CloneAggregation(aggregation_type_, instrument_descriptor_,
*observed);
}
last_observed_hash_map_->Set(attributes, std::move(observed));

AccumulateDelta(FilterAttributes(attributes), *delta);
}

/**
* Adds one series' delta to its output point. Merging, not overwriting, is what lets the series
* a view collapses together and the observations from different callbacks all contribute.
* Over-the-limit attribute sets resolve to the same otel.metric.overflow entry, which merges
* too.
*/
void AccumulateDelta(const MetricAttributes &attributes, const Aggregation &delta)
{
auto merged =
delta_hash_map_->GetOrSetDefault(attributes, create_default_aggregation_)->Merge(delta);
delta_hash_map_->Set(attributes, std::move(merged));
}

/**
* Differences this round's absolute observations and starts a fresh round, for everything but
* monotonic sums. An attribute set missing from this round keeps its baseline and contributes
* no delta; suppressing the stale output is the temporal storage's job.
*/
void BuildDeltaFromRound() noexcept
{
round_hash_map_->GetAllEntries([this](const MetricAttributes &attributes,
Aggregation &aggregation) {
auto observed = DefaultAggregation::CloneAggregation(aggregation_type_,
instrument_descriptor_, aggregation);
auto previous = last_observed_hash_map_->Get(attributes);
if (previous)
{
delta_hash_map_->Set(attributes, previous->Diff(*observed));
}
else
{
delta_hash_map_->Set(attributes, DefaultAggregation::CloneAggregation(
aggregation_type_, instrument_descriptor_, *observed));
}
last_observed_hash_map_->Set(attributes, std::move(observed));
return true;
});

round_hash_map_ = std::make_unique<AttributesHashMap>(aggregation_config_->cardinality_limit_);
}

/**
* Returns a copy of the observed attributes with the attributes dropped by the view removed.
*/
MetricAttributes FilterAttributes(const MetricAttributes &attributes) const
{
MetricAttributes filtered(attributes);
if (!attributes_processor_)
{
return filtered;
}

bool dropped = false;
for (auto iter = filtered.begin(); iter != filtered.end();)
{
if (attributes_processor_->isPresent(iter->first))
{
++iter;
}
else
{
iter = filtered.erase(iter);
dropped = true;
}
}

if (dropped)
{
filtered.UpdateHash();
}
return filtered;
}

/**
* Whether this storage aggregates into a monotonic sum, matching what
* DefaultAggregation::CreateAggregation() would create.
*/
static bool IsMonotonicSum(AggregationType aggregation_type,
const InstrumentDescriptor &instrument_descriptor) noexcept
{
bool is_monotonic = true;
if (aggregation_type == AggregationType::kDefault)
{
const AggregationType resolved =
DefaultAggregation::GetDefaultAggregationType(instrument_descriptor.type_, is_monotonic);
return resolved == AggregationType::kSum && is_monotonic;
}
if (aggregation_type != AggregationType::kSum)
{
return false;
}
return instrument_descriptor.type_ != InstrumentType::kUpDownCounter &&
instrument_descriptor.type_ != InstrumentType::kObservableUpDownCounter &&
instrument_descriptor.type_ != InstrumentType::kHistogram;
}

InstrumentDescriptor instrument_descriptor_;
AggregationType aggregation_type_;
const AggregationConfig *aggregation_config_;
std::unique_ptr<AttributesHashMap> cumulative_hash_map_;
std::shared_ptr<const AttributesProcessor> attributes_processor_;
bool is_monotonic_sum_;
std::function<std::unique_ptr<Aggregation>()> create_default_aggregation_;

std::unique_ptr<AttributesHashMap> last_observed_hash_map_;
std::unique_ptr<AttributesHashMap> delta_hash_map_;
std::unique_ptr<AttributesHashMap> round_hash_map_;

std::mutex hashmap_lock_;
#ifdef ENABLE_METRICS_EXEMPLAR_PREVIEW
ExemplarFilterType exemplar_filter_type_;
Expand Down
2 changes: 1 addition & 1 deletion sdk/src/metrics/meter.cc
Original file line number Diff line number Diff line change
Expand Up @@ -615,7 +615,7 @@ std::unique_ptr<AsyncWritableMetricStorage> Meter::RegisterAsyncMetricStorage(
{
WarnOnDuplicateInstrument(GetInstrumentationScope(), storage_registry_, view_instr_desc);
async_storage = std::shared_ptr<AsyncMetricStorage>(new AsyncMetricStorage(
view_instr_desc, view.GetAggregationType(),
view_instr_desc, view.GetAggregationType(), view.GetAttributesProcessor(),
#ifdef ENABLE_METRICS_EXEMPLAR_PREVIEW
exemplar_filter_type,
GetExemplarReservoir(view.GetAggregationType(), view.GetAggregationConfig(),
Expand Down
Loading
Loading