Skip to content

Commit e37160d

Browse files
authored
fix: use one timestamp for OTEL metric batches (#1292)
Quick fix to prevent batch chunks having out of order timestamps, which messes up certain strict OTEL ingest pipelines like described in #1289. This makes every chunk have the same timestamp, so if duplicate data is sent in different chunks, particularly metadata such as [target_info](https://prometheus.io/docs/guides/opentelemetry/#including-resource-attributes-at-query-time), it dedupes cleanly.
1 parent 6039f31 commit e37160d

2 files changed

Lines changed: 14 additions & 13 deletions

File tree

pgdog/src/stats/otel.rs

Lines changed: 10 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -31,7 +31,7 @@ struct CounterKey {
3131
static PREV_COUNTERS: Lazy<Mutex<HashMap<CounterKey, f64>>> =
3232
Lazy::new(|| Mutex::new(HashMap::new()));
3333

34-
fn now_nanos() -> String {
34+
pub fn now_nanos() -> String {
3535
SystemTime::now()
3636
.duration_since(UNIX_EPOCH)
3737
.unwrap_or_default()
@@ -203,8 +203,7 @@ fn measurement_to_f64(m: &MeasurementType) -> f64 {
203203
}
204204

205205
/// Build an `ExportMetricsServiceRequest` from a collection of `Metric` objects.
206-
pub fn build_request(metrics: &[&Metric]) -> ExportMetricsServiceRequest {
207-
let now = now_nanos();
206+
pub fn build_request(metrics: &[&Metric], now: &str) -> ExportMetricsServiceRequest {
208207
let config = crate::config::config();
209208
let namespace = config
210209
.config
@@ -271,7 +270,7 @@ pub fn build_request(metrics: &[&Metric]) -> ExportMetricsServiceRequest {
271270

272271
Some(NumberDataPoint {
273272
start_time_unix_nano: None,
274-
time_unix_nano: now.clone(),
273+
time_unix_nano: now.to_owned(),
275274
as_double,
276275
attributes,
277276
})
@@ -340,7 +339,7 @@ mod test {
340339
metric_type: None,
341340
});
342341

343-
let request = build_request(&[&metric]);
342+
let request = build_request(&[&metric], &now_nanos());
344343
let json = serde_json::to_string_pretty(&request).expect("serialize");
345344

346345
assert!(json.contains("\"gauge\""));
@@ -365,7 +364,7 @@ mod test {
365364
metric_type: Some("counter".into()),
366365
});
367366

368-
let request = build_request(&[&metric]);
367+
let request = build_request(&[&metric], &now_nanos());
369368
let json = serde_json::to_string(&request).expect("serialize");
370369

371370
assert!(json.contains("\"sum\""));
@@ -396,7 +395,7 @@ mod test {
396395
metric_type: None,
397396
});
398397

399-
let request = build_request(&[&metric]);
398+
let request = build_request(&[&metric], &now_nanos());
400399
let scope = &request.resource_metrics[0].scope_metrics[0].scope;
401400
assert_eq!(scope.name, "pgdog");
402401

@@ -428,7 +427,7 @@ mod test {
428427
metric_type: None,
429428
});
430429

431-
let request = build_request(&[&metric]);
430+
let request = build_request(&[&metric], &now_nanos());
432431
let scope = &request.resource_metrics[0].scope_metrics[0];
433432
let otel_metric = &scope.metrics[0];
434433

@@ -458,7 +457,7 @@ mod test {
458457
metric_type: Some("counter".into()),
459458
});
460459

461-
let request = build_request(&[&metric]);
460+
let request = build_request(&[&metric], &now_nanos());
462461
let sum = &request.resource_metrics[0].scope_metrics[0].metrics[0]
463462
.sum
464463
.as_ref()
@@ -481,7 +480,7 @@ mod test {
481480
metric_type: None,
482481
});
483482

484-
let request = build_request(&[&metric]);
483+
let request = build_request(&[&metric], &now_nanos());
485484
let gauge = &request.resource_metrics[0].scope_metrics[0].metrics[0]
486485
.gauge
487486
.as_ref()
@@ -505,7 +504,7 @@ mod test {
505504
metric_type: None,
506505
});
507506

508-
let request = build_request(&[&metric]);
507+
let request = build_request(&[&metric], &now_nanos());
509508
let resource = &request.resource_metrics[0].resource;
510509

511510
let svc = resource

pgdog/src/stats/otel_exporter.rs

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -54,11 +54,13 @@ pub async fn run() {
5454
all.extend(listeners.iter());
5555
all.extend(query_cache.iter());
5656

57+
let now = otel::now_nanos();
58+
5759
// Send batches in parallel to stay under the 512 KB payload limit.
5860
let futs: Vec<_> = all
5961
.chunks(BATCH_SIZE)
6062
.filter_map(|chunk| {
61-
let request = otel::build_request(chunk);
63+
let request = otel::build_request(chunk, &now);
6264
let body = match serde_json::to_vec(&request) {
6365
Ok(b) => b,
6466
Err(err) => {
@@ -124,7 +126,7 @@ mod test {
124126
metric_type: None,
125127
});
126128

127-
let request = otel::build_request(&[&metric]);
129+
let request = otel::build_request(&[&metric], &otel::now_nanos());
128130
let body = serde_json::to_vec(&request).expect("serialize");
129131

130132
// Verify the output is valid JSON by parsing it back.

0 commit comments

Comments
 (0)