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
3 changes: 2 additions & 1 deletion tpu_sync/telemetry/BUILD
Original file line number Diff line number Diff line change
Expand Up @@ -118,10 +118,10 @@ cc_library(
hdrs = ["buffered_metrics_exporter.h"],
deps = [
":exporter_util",
":label_util",
":metrics_backend",
"@com_google_absl//absl/base:core_headers",
"@com_google_absl//absl/container:flat_hash_map",
"@com_google_absl//absl/container:inlined_vector",
"@com_google_absl//absl/log",
"@com_google_absl//absl/strings",
"@com_google_absl//absl/synchronization",
Expand All @@ -134,6 +134,7 @@ cc_test(
srcs = ["buffered_metrics_exporter_test.cc"],
deps = [
":buffered_metrics_exporter",
":label_util",
":metrics_backend",
"@com_google_absl//absl/strings",
"@com_google_absl//absl/types:span",
Expand Down
83 changes: 12 additions & 71 deletions tpu_sync/telemetry/buffered_metrics_exporter.cc
Original file line number Diff line number Diff line change
Expand Up @@ -14,15 +14,13 @@

#include "tpu_sync/telemetry/buffered_metrics_exporter.h"

#include <algorithm>
#include <cstdint>
#include <map>
#include <memory>
#include <string>
#include <utility>
#include <vector>

#include "absl/container/inlined_vector.h"
#include "absl/strings/str_cat.h"
#include "absl/strings/string_view.h"
#include "absl/types/span.h"
Expand All @@ -31,74 +29,18 @@

namespace tpu_raiden::telemetry {

namespace {

std::string EscapeLabelValue(absl::string_view value) {
std::string escaped;
escaped.reserve(value.size());
for (char c : value) {
switch (c) {
case '\\':
escaped.push_back('\\');
escaped.push_back('\\');
break;
case '"':
escaped.push_back('\\');
escaped.push_back('"');
break;
case '\n':
escaped.push_back('\\');
escaped.push_back('n');
break;
default:
escaped.push_back(c);
break;
}
}
return escaped;
}

} // namespace

std::string FormatCanonicalLabels(LabelSpan labels) {
if (labels.empty()) {
return "";
}
absl::InlinedVector<MetricLabel, kDefaultInlinedLabelCapacity> sorted_labels(
labels.begin(), labels.end());
std::sort(sorted_labels.begin(), sorted_labels.end());

std::string out;
out.push_back('{');
bool first = true;
for (const auto& [key, value] : sorted_labels) {
if (!first) {
out.push_back(',');
}
first = false;
out.append(key);
out.push_back('=');
out.push_back('"');
out.append(EscapeLabelValue(value));
out.push_back('"');
}
out.push_back('}');
return out;
}

BufferedMetricsExporter::BufferedMetricsExporter(
absl::Span<const MetricMetadata> metrics) {
for (const auto& meta : metrics) {
switch (meta.type) {
case MetricType::kCounter:
counters_.emplace(
meta.name,
std::make_unique<
MetricFamilyBuffer<LockFreeCounterAccumulator>>());
std::make_unique<MetricFamilyBuffer<LockFreeCounterAccumulator>>());
break;
case MetricType::kGauge:
gauges_.emplace(
meta.name, std::make_unique<MetricFamilyBuffer<QueueBuffer<>>>());
gauges_.emplace(meta.name,
std::make_unique<MetricFamilyBuffer<QueueBuffer<>>>());
break;
case MetricType::kHistogram:
histograms_.emplace(
Expand Down Expand Up @@ -148,16 +90,15 @@ BufferedMetricsExporter::GetAndResetMetricSamples() {
std::map<std::string, std::vector<double>> result;

for (const auto& [name, family_buffer] : counters_) {
family_buffer->ForEachAccumulator(
[&](absl::string_view canonical_labels,
LockFreeCounterAccumulator* counter) {
uint64_t delta = counter->ExchangeAndReset();
if (delta > 0) {
const std::string full_name =
absl::StrCat(kPrometheusMetricPrefix, name, canonical_labels);
result[full_name].push_back(static_cast<double>(delta));
}
});
family_buffer->ForEachAccumulator([&](absl::string_view canonical_labels,
LockFreeCounterAccumulator* counter) {
uint64_t delta = counter->ExchangeAndReset();
if (delta > 0) {
const std::string full_name =
absl::StrCat(kPrometheusMetricPrefix, name, canonical_labels);
result[full_name].push_back(static_cast<double>(delta));
}
});
}

for (const auto& [name, family_buffer] : gauges_) {
Expand Down
41 changes: 20 additions & 21 deletions tpu_sync/telemetry/buffered_metrics_exporter.h
Original file line number Diff line number Diff line change
Expand Up @@ -21,7 +21,6 @@
#include <map>
#include <memory>
#include <string>
#include <utility>
#include <vector>

#include "absl/base/thread_annotations.h"
Expand All @@ -30,6 +29,7 @@
#include "absl/strings/string_view.h"
#include "absl/synchronization/mutex.h"
#include "absl/types/span.h"
#include "tpu_sync/telemetry/label_util.h"
#include "tpu_sync/telemetry/metrics_backend.h"

namespace tpu_raiden::telemetry {
Expand All @@ -51,9 +51,8 @@ class QueueBuffer {
if (samples_.size() < max_capacity_) {
samples_.push_back(val);
} else {
LOG_EVERY_N_SEC(WARNING, 5)
<< "QueueBuffer capacity (" << max_capacity_
<< ") exceeded; dropping metric sample.";
LOG_EVERY_N_SEC(WARNING, 5) << "QueueBuffer capacity (" << max_capacity_
<< ") exceeded; dropping metric sample.";
}
}

Expand Down Expand Up @@ -90,10 +89,6 @@ class LockFreeCounterAccumulator {
};

inline constexpr size_t kMaxLabeledSeries = 512;
// Formats canonical Prometheus label string for in-memory series
// identification. Returns "{key1=\"val1\",key2=\"val2\"}" sorted by key, with
// Prometheus character escaping. Returns "" if labels are empty.
std::string FormatCanonicalLabels(LabelSpan labels);

// Encapsulates all time-series buffers for a single metric family.
template <typename AccumulatorType>
Expand All @@ -113,27 +108,33 @@ class MetricFamilyBuffer {
if (labels.empty()) {
return &unlabeled_;
}
std::string canonical_labels = FormatCanonicalLabels(labels);
PrometheusLabelView formatted(labels);
const absl::string_view lookup_key = formatted.view();
if (lookup_key.empty()) {
LOG_EVERY_N_SEC(WARNING, 5)
<< "Failed to format metric labels; dropping metric series.";
return nullptr;
}

{
absl::ReaderMutexLock lock(labeled_mu_);
if (auto it = labeled_.find(canonical_labels); it != labeled_.end()) {
if (auto it = labeled_.find(lookup_key); it != labeled_.end()) {
return it->second.get();
}
}

absl::MutexLock lock(labeled_mu_);
auto it = labeled_.find(canonical_labels);
if (it != labeled_.end()) {
if (auto it = labeled_.find(lookup_key); it != labeled_.end()) {
return it->second.get();
}
if (labeled_.size() >= kMaxLabeledSeries) {
LOG_EVERY_N_SEC(WARNING, 5)
<< "Max labeled series capacity (" << kMaxLabeledSeries
<< ") exceeded; dropping metric series for labels: "
<< canonical_labels;
<< ") exceeded; dropping metric series for labels: " << lookup_key;
return nullptr;
}
auto [insert_it, _] = labeled_.try_emplace(
std::move(canonical_labels), std::make_unique<AccumulatorType>());
formatted.ToOwned(), std::make_unique<AccumulatorType>());
return insert_it->second.get();
}

Expand Down Expand Up @@ -189,13 +190,11 @@ class BufferedMetricsExporter : public MetricsBackend {
std::string,
std::unique_ptr<MetricFamilyBuffer<LockFreeCounterAccumulator>>>
counters_;
absl::flat_hash_map<
std::string,
std::unique_ptr<MetricFamilyBuffer<QueueBuffer<>>>>
absl::flat_hash_map<std::string,
std::unique_ptr<MetricFamilyBuffer<QueueBuffer<>>>>
gauges_;
absl::flat_hash_map<
std::string,
std::unique_ptr<MetricFamilyBuffer<QueueBuffer<>>>>
absl::flat_hash_map<std::string,
std::unique_ptr<MetricFamilyBuffer<QueueBuffer<>>>>
histograms_;
};

Expand Down
Loading
Loading