Skip to content

Commit 29d1b20

Browse files
committed
agents: implement custom OTLP metrics serializer
Introduce custom OTLP metrics serialization in gRPC agent, replacing SDK with BatchedMetricData and direct protobuf generation. Provides per-data-point timestamps, custom batching, and improved performance for high-frequency metrics. Add debug logging, fix thread ID type mismatches, implement robust batch processing for partial metrics, and update test expectations to resolve inconsistent processing and hangs.
1 parent 98c413a commit 29d1b20

12 files changed

Lines changed: 350 additions & 69 deletions

agents/grpc/src/grpc_agent.cc

Lines changed: 11 additions & 22 deletions
Original file line numberDiff line numberDiff line change
@@ -265,9 +265,15 @@ void PopulateMetricsEvent(grpcagent::MetricsEvent* metrics_event,
265265
otlp::fill_env_metrics(metrics, env_metrics_stor, false);
266266
}
267267

268+
auto batched = metrics.DumpMetricsAndReset();
269+
std::vector<opentelemetry::sdk::metrics::MetricData> metric_data;
270+
for (const auto& bm : batched) {
271+
metric_data.push_back(otlp::BatchedMetricToMetricData(bm));
272+
}
273+
268274
data.scope_metric_data_ =
269275
std::vector<ScopeMetrics>{{ otlp::GetScope(),
270-
metrics.DumpMetricsAndReset() }};
276+
std::move(metric_data) }};
271277
OtlpMetricUtils::PopulateResourceMetrics(
272278
data, metrics_event->mutable_body()->mutable_resource_metrics()->Add());
273279
}
@@ -946,11 +952,7 @@ void GrpcAgent::env_deletion_cb_(SharedEnvInst envinst,
946952
if (agent->thr_metrics_batch_.ShouldFlush()) {
947953
auto metric_data = agent->thr_metrics_batch_.DumpMetricsAndReset();
948954
if (!metric_data.empty()) {
949-
ResourceMetrics data;
950-
data.resource_ = otlp::GetResource();
951-
data.scope_metric_data_ =
952-
std::vector<ScopeMetrics>{{ otlp::GetScope(), std::move(metric_data) }};
953-
agent->metrics_exporter_->enqueue(std::move(data));
955+
agent->metrics_exporter_->enqueue(std::move(metric_data));
954956
}
955957
}
956958
}
@@ -969,24 +971,15 @@ void GrpcAgent::on_metrics_timer() {
969971
if (metrics_paused_) {
970972
auto metric_data = proc_metrics_batch_.DumpMetricsAndReset();
971973
if (!metric_data.empty()) {
972-
ResourceMetrics data;
973-
data.resource_ = otlp::GetResource();
974-
data.scope_metric_data_ =
975-
std::vector<ScopeMetrics>{{ otlp::GetScope(), std::move(metric_data) }};
976-
metrics_exporter_->enqueue(std::move(data));
974+
metrics_exporter_->enqueue(std::move(metric_data));
977975
}
978976

979977
metric_data = thr_metrics_batch_.DumpMetricsAndReset();
980978
if (!metric_data.empty()) {
981-
ResourceMetrics data;
982-
data.resource_ = otlp::GetResource();
983-
data.scope_metric_data_ =
984-
std::vector<ScopeMetrics>{{ otlp::GetScope(), std::move(metric_data) }};
985-
metrics_exporter_->enqueue(std::move(data));
979+
metrics_exporter_->enqueue(std::move(metric_data));
986980
}
987981

988982
metrics_exporter_->flush();
989-
return;
990983
}
991984

992985
got_proc_metrics();
@@ -1459,11 +1452,7 @@ void GrpcAgent::got_proc_metrics() {
14591452
if (proc_metrics_batch_.ShouldFlush()) {
14601453
auto metric_data = proc_metrics_batch_.DumpMetricsAndReset();
14611454
if (!metric_data.empty()) {
1462-
ResourceMetrics data;
1463-
data.resource_ = otlp::GetResource();
1464-
data.scope_metric_data_ =
1465-
std::vector<ScopeMetrics>{{ otlp::GetScope(), std::move(metric_data) }};
1466-
metrics_exporter_->enqueue(std::move(data));
1455+
metrics_exporter_->enqueue(std::move(metric_data));
14671456
}
14681457
}
14691458
}

agents/grpc/src/grpc_metrics_exporter.cc

Lines changed: 5 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,7 @@
11
#include "grpc_metrics_exporter.h"
22
#include "grpc_client.h"
33

4+
#include "../../otlp/src/otlp_common.h"
45
#include "opentelemetry/exporters/otlp/otlp_grpc_client.h"
56
#include "opentelemetry/exporters/otlp/otlp_metric_utils.h"
67

@@ -52,7 +53,8 @@ void GrpcMetricsExporter::init() {
5253
GrpcMetricsExporter::~GrpcMetricsExporter() {
5354
}
5455

55-
void GrpcMetricsExporter::enqueue(ResourceMetrics&& metrics) {
56+
void GrpcMetricsExporter::enqueue(
57+
std::vector<otlp::BatchedMetricData>&& metrics) {
5658
metrics_q_.push(std::move(metrics));
5759
if (!in_flight_) {
5860
export_current();
@@ -78,7 +80,8 @@ void GrpcMetricsExporter::export_current() {
7880

7981
auto* request =
8082
google::protobuf::Arena::Create<ExportMetricsServiceRequest>(arena.get());
81-
OtlpMetricUtils::PopulateRequest(data, request);
83+
84+
otlp::PopulateRequest(data, otlp::GetResource(), otlp::GetScope(), request);
8285

8386
auto context = OtlpGrpcClient::MakeClientContext(options_);
8487
::grpc::Status immediate =

agents/grpc/src/grpc_metrics_exporter.h

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,7 @@
55
#include <memory>
66
#include <string>
77

8+
#include "../../otlp/src/batched_metric_data.h"
89
#include "nsolid/async_ts_queue.h"
910
#include "nsolid/nsolid_util.h"
1011
#include "opentelemetry/exporters/otlp/otlp_grpc_client.h"
@@ -45,7 +46,7 @@ class GrpcMetricsExporter:
4546

4647
void init();
4748

48-
void enqueue(opentelemetry::sdk::metrics::ResourceMetrics&& metrics);
49+
void enqueue(std::vector<otlp::BatchedMetricData>&& metrics);
4950

5051
void flush();
5152

@@ -63,7 +64,7 @@ class GrpcMetricsExporter:
6364
opentelemetry::v1::exporter::otlp::OtlpGrpcMetricExporterOptions options_;
6465
std::shared_ptr<opentelemetry::v1::exporter::otlp::OtlpGrpcClient> client_;
6566
std::unique_ptr<MetricsServiceStub> metrics_service_stub_;
66-
utils::RingBuffer<opentelemetry::sdk::metrics::ResourceMetrics> metrics_q_;
67+
utils::RingBuffer<std::vector<otlp::BatchedMetricData>> metrics_q_;
6768
bool in_flight_;
6869
std::shared_ptr<AsyncTSQueue<::grpc::Status>> metrics_completion_q_;
6970
};
Lines changed: 51 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,51 @@
1+
#ifndef AGENTS_OTLP_SRC_BATCHED_METRIC_DATA_H_
2+
#define AGENTS_OTLP_SRC_BATCHED_METRIC_DATA_H_
3+
4+
#include <vector>
5+
6+
#include "opentelemetry/nostd/variant.h"
7+
#include "opentelemetry/sdk/common/attribute_utils.h"
8+
#include "opentelemetry/sdk/metrics/data/point_data.h"
9+
#include "opentelemetry/sdk/metrics/instruments.h"
10+
#include "opentelemetry/version.h"
11+
12+
namespace node {
13+
namespace nsolid {
14+
namespace otlp {
15+
16+
using opentelemetry::sdk::metrics::SumPointData;
17+
using opentelemetry::sdk::metrics::HistogramPointData;
18+
using opentelemetry::sdk::metrics::Base2ExponentialHistogramPointData;
19+
using opentelemetry::sdk::metrics::LastValuePointData;
20+
using opentelemetry::sdk::metrics::DropPointData;
21+
using opentelemetry::sdk::metrics::SummaryPointData;
22+
using opentelemetry::sdk::metrics::AggregationTemporality;
23+
using opentelemetry::sdk::metrics::InstrumentDescriptor;
24+
using PointAttributes = opentelemetry::sdk::common::OrderedAttributeMap;
25+
using PointType =
26+
opentelemetry::nostd::variant<SumPointData,
27+
HistogramPointData,
28+
Base2ExponentialHistogramPointData,
29+
LastValuePointData,
30+
DropPointData,
31+
SummaryPointData>;
32+
33+
struct TimedPointDataAttributes {
34+
opentelemetry::common::SystemTimestamp start_ts;
35+
opentelemetry::common::SystemTimestamp end_ts;
36+
PointAttributes attributes;
37+
PointType point_data;
38+
};
39+
40+
class BatchedMetricData {
41+
public:
42+
InstrumentDescriptor instrument_descriptor;
43+
AggregationTemporality aggregation_temporality;
44+
std::vector<TimedPointDataAttributes> point_data_attr_;
45+
};
46+
47+
} // namespace otlp
48+
} // namespace nsolid
49+
} // namespace node
50+
51+
#endif // AGENTS_OTLP_SRC_BATCHED_METRIC_DATA_H_

0 commit comments

Comments
 (0)