From 477292b2a7df688f230a794ac7720d85f6acb2a3 Mon Sep 17 00:00:00 2001 From: Jesse Szwedko Date: Fri, 2 Oct 2026 15:17:05 -0400 Subject: [PATCH] feat(agent-data-plane): add aggregator flush count and timing telemetry Co-Authored-By: Claude Opus 5.5 (1M context) --- bin/agent-data-plane/src/state/metrics/mod.rs | 91 +++++++++++++++++++ .../src/state/metrics/rules/aggregation.rs | 72 ++++++++++++++- .../src/encoders/buffered_incremental/mod.rs | 38 ++++++-- .../buffered_incremental/telemetry.rs | 38 +++++++- .../src/transforms/aggregate/mod.rs | 61 ++++++++++++- .../src/transforms/aggregate/telemetry.rs | 31 ++++++- ...ator-flush-telemetry-943c788687171dd3.yaml | 6 ++ 7 files changed, 323 insertions(+), 14 deletions(-) create mode 100644 releasenotes/notes/rar-aggregator-flush-telemetry-943c788687171dd3.yaml diff --git a/bin/agent-data-plane/src/state/metrics/mod.rs b/bin/agent-data-plane/src/state/metrics/mod.rs index aad9e19c491..930aa9089c1 100644 --- a/bin/agent-data-plane/src/state/metrics/mod.rs +++ b/bin/agent-data-plane/src/state/metrics/mod.rs @@ -571,6 +571,14 @@ mod tests { find("aggregator.number_of_flush"), Some("Number of flushes done by the aggregator") ); + assert_eq!( + find("aggregator.flush_count"), + Some("Number of items handled by the last flush, by flush type") + ); + assert_eq!( + find("aggregator.flush_time"), + Some("Duration in nanoseconds of the last flush, by flush type") + ); assert_eq!(find("filterlist.size"), Some("Metric filter list size")); assert_eq!( find("filterlist.updates"), @@ -693,6 +701,89 @@ mod tests { assert!(output.contains("# TYPE aggregator__number_of_flush counter")); } + #[test] + fn render_rar_telemetry_remaps_aggregator_flush_telemetry() { + let metrics = vec![ + Event::Metric(Metric::gauge( + Context::from_static_parts( + "adp.aggregate_last_flush_count", + &["component_id:dsd_agg", "data_type:series"], + ), + 11.0, + )), + Event::Metric(Metric::gauge( + Context::from_static_parts( + "adp.aggregate_last_flush_count", + &["component_id:dsd_agg", "data_type:sketches"], + ), + 12.0, + )), + Event::Metric(Metric::gauge( + Context::from_static_parts( + "adp.aggregate_last_flush_duration_nanoseconds", + &["component_id:dsd_agg"], + ), + 1000.0, + )), + Event::Metric(Metric::counter( + Context::from_static_parts("adp.encoder_flushed_events_total", &["component_id:dd_events_encode"]), + 5.0, + )), + Event::Metric(Metric::counter( + Context::from_static_parts( + "adp.encoder_flushed_events_total", + &["component_id:dd_service_checks_encode"], + ), + 6.0, + )), + Event::Metric(Metric::gauge( + Context::from_static_parts("adp.encoder_last_flush_events", &["component_id:dd_events_encode"]), + 2.0, + )), + Event::Metric(Metric::gauge( + Context::from_static_parts( + "adp.encoder_last_flush_events", + &["component_id:dd_service_checks_encode"], + ), + 3.0, + )), + Event::Metric(Metric::gauge( + Context::from_static_parts( + "adp.encoder_last_flush_duration_nanoseconds", + &["component_id:dd_events_encode"], + ), + 2000.0, + )), + Event::Metric(Metric::gauge( + Context::from_static_parts( + "adp.encoder_last_flush_duration_nanoseconds", + &["component_id:dd_service_checks_encode"], + ), + 3000.0, + )), + // The logs encoder shares the same encoder telemetry, but has no Core Agent aggregator counterpart. + Event::Metric(Metric::counter( + Context::from_static_parts("adp.encoder_flushed_events_total", &["component_id:dd_logs_encode"]), + 99.0, + )), + ]; + + let output = render_with(get_datadog_agent_remappings(), metrics); + + assert!(output.contains("aggregator__flush_count{flush_type=\"series\"} 11")); + assert!(output.contains("aggregator__flush_count{flush_type=\"sketches\"} 12")); + assert!(output.contains("aggregator__flush_count{flush_type=\"events\"} 2")); + assert!(output.contains("aggregator__flush_count{flush_type=\"service_checks\"} 3")); + assert!(output.contains("aggregator__flush_time{flush_type=\"main\"} 1000")); + assert!(output.contains("aggregator__flush_time{flush_type=\"event\"} 2000")); + assert!(output.contains("aggregator__flush_time{flush_type=\"service_check\"} 3000")); + assert!(output.contains("aggregator__flush{data_type=\"events\"} 5")); + assert!(output.contains("aggregator__flush{data_type=\"service_checks\"} 6")); + assert!(!output.contains(" 99")); + assert!(output.contains("# TYPE aggregator__flush_count gauge")); + assert!(output.contains("# TYPE aggregator__flush_time gauge")); + } + #[test] fn rar_telemetry_remaps_supported_dogstatsd_client_byte_telemetry_as_counters() { let metrics = [ diff --git a/bin/agent-data-plane/src/state/metrics/rules/aggregation.rs b/bin/agent-data-plane/src/state/metrics/rules/aggregation.rs index 701f3778571..6dbb8a0ce85 100644 --- a/bin/agent-data-plane/src/state/metrics/rules/aggregation.rs +++ b/bin/agent-data-plane/src/state/metrics/rules/aggregation.rs @@ -1,6 +1,12 @@ use super::RemapperRule; const NO_AGG_SPLIT_COMPONENT_TAG: &str = "component_id:dsd_no_agg_split"; +const EVENTS_ENCODER_COMPONENT_TAG: &str = "component_id:dd_events_encode"; +const SERVICE_CHECKS_ENCODER_COMPONENT_TAG: &str = "component_id:dd_service_checks_encode"; + +const FLUSH_HELP_TEXT: &str = "Number of metrics/service checks/events flushed"; +const FLUSH_COUNT_HELP_TEXT: &str = "Number of items handled by the last flush, by flush type"; +const FLUSH_TIME_HELP_TEXT: &str = "Duration in nanoseconds of the last flush, by flush type"; pub fn get_aggregation_remappings() -> Vec { vec![ @@ -34,14 +40,14 @@ pub fn get_aggregation_remappings() -> Vec { // Events and service checks don't pass through an aggregator in ADP, so their encoders stand in for it. RemapperRule::by_name_and_tags( "adp.component_events_received_total", - &["component_id:dd_events_encode"], + &[EVENTS_ENCODER_COMPONENT_TAG], "aggregator.processed", ) .with_additional_tags(["data_type:events"]) .with_help_text("Amount of metrics/services_checks/events processed by the aggregator"), RemapperRule::by_name_and_tags( "adp.component_events_received_total", - &["component_id:dd_service_checks_encode"], + &[SERVICE_CHECKS_ENCODER_COMPONENT_TAG], "aggregator.processed", ) .with_additional_tags(["data_type:service_checks"]) @@ -65,12 +71,72 @@ pub fn get_aggregation_remappings() -> Vec { "aggregator.flush", ) .with_original_tags(["data_type"]) - .with_help_text("Number of metrics/service checks/events flushed"), + .with_help_text(FLUSH_HELP_TEXT), RemapperRule::by_name_and_tags( "adp.aggregate_flushes_total", &["component_id:dsd_agg"], "aggregator.number_of_flush", ) .with_help_text("Number of flushes done by the aggregator"), + RemapperRule::by_name_and_tags( + "adp.aggregate_last_flush_count", + &["component_id:dsd_agg"], + "aggregator.flush_count", + ) + .with_remapped_tags([("data_type", "flush_type")]) + .with_help_text(FLUSH_COUNT_HELP_TEXT), + // There's no ADP equivalent of the Core Agent's separate `metric_sketch` and `checks_metric_sample` flush + // timings, which include serialization: ADP encodes metrics in a separate component after the aggregator flush. + RemapperRule::by_name_and_tags( + "adp.aggregate_last_flush_duration_nanoseconds", + &["component_id:dsd_agg"], + "aggregator.flush_time", + ) + .with_additional_tags(["flush_type:main"]) + .with_help_text(FLUSH_TIME_HELP_TEXT), + // As with `aggregator.processed`, the events and service checks encoders stand in for the aggregator's + // event and service check flushes. + RemapperRule::by_name_and_tags( + "adp.encoder_flushed_events_total", + &[EVENTS_ENCODER_COMPONENT_TAG], + "aggregator.flush", + ) + .with_additional_tags(["data_type:events"]) + .with_help_text(FLUSH_HELP_TEXT), + RemapperRule::by_name_and_tags( + "adp.encoder_flushed_events_total", + &[SERVICE_CHECKS_ENCODER_COMPONENT_TAG], + "aggregator.flush", + ) + .with_additional_tags(["data_type:service_checks"]) + .with_help_text(FLUSH_HELP_TEXT), + RemapperRule::by_name_and_tags( + "adp.encoder_last_flush_events", + &[EVENTS_ENCODER_COMPONENT_TAG], + "aggregator.flush_count", + ) + .with_additional_tags(["flush_type:events"]) + .with_help_text(FLUSH_COUNT_HELP_TEXT), + RemapperRule::by_name_and_tags( + "adp.encoder_last_flush_events", + &[SERVICE_CHECKS_ENCODER_COMPONENT_TAG], + "aggregator.flush_count", + ) + .with_additional_tags(["flush_type:service_checks"]) + .with_help_text(FLUSH_COUNT_HELP_TEXT), + RemapperRule::by_name_and_tags( + "adp.encoder_last_flush_duration_nanoseconds", + &[EVENTS_ENCODER_COMPONENT_TAG], + "aggregator.flush_time", + ) + .with_additional_tags(["flush_type:event"]) + .with_help_text(FLUSH_TIME_HELP_TEXT), + RemapperRule::by_name_and_tags( + "adp.encoder_last_flush_duration_nanoseconds", + &[SERVICE_CHECKS_ENCODER_COMPONENT_TAG], + "aggregator.flush_time", + ) + .with_additional_tags(["flush_type:service_check"]) + .with_help_text(FLUSH_TIME_HELP_TEXT), ] } diff --git a/lib/saluki-components/src/encoders/buffered_incremental/mod.rs b/lib/saluki-components/src/encoders/buffered_incremental/mod.rs index 26714789e1d..294903b9197 100644 --- a/lib/saluki-components/src/encoders/buffered_incremental/mod.rs +++ b/lib/saluki-components/src/encoders/buffered_incremental/mod.rs @@ -1,4 +1,4 @@ -use std::time::Duration; +use std::time::{Duration, Instant}; use agent_data_plane_config::defaults::DEFAULT_ENCODER_FLUSH_TIMEOUT; use async_trait::async_trait; @@ -9,6 +9,7 @@ use saluki_core::{ components::{encoders::*, BuildContext}, data_model::{event::EventType, payload::PayloadType}, observability::ComponentMetricsExt, + topology::PayloadsDispatcher, }; use saluki_error::GenericError; use saluki_metrics::MetricsBuilder; @@ -173,6 +174,7 @@ where health.mark_ready(); let mut pending_flush = false; + let mut pending_events = 0; let pending_flush_timeout = sleep(flush_timeout); pin!(pending_flush_timeout); @@ -192,19 +194,22 @@ where // If we're informed that we need to flush, we'll hold on to this event before triggering a flush and then // retry processing it after flushing. let event_to_retry = match encoder.process_event(event).await? { - ProcessResult::Continue => continue, + ProcessResult::Continue => { + pending_events += 1; + continue; + } ProcessResult::FlushRequired(event) => event, }; // Flush the encoder, waiting any payloads it has generated. - encoder.flush(context.dispatcher()).await?; + flush_encoder(&mut encoder, context.dispatcher(), &telemetry, &mut pending_events).await?; // Now try to process the event again. // // If this fails, then we drop the event because it's a logical bug to not be able to encode an event after // flushing, and we don't want to get stuck in an infinite loop. match encoder.process_event(event_to_retry).await? { - ProcessResult::Continue => {}, + ProcessResult::Continue => pending_events += 1, ProcessResult::FlushRequired(_) => { error!("Failed to process event after flushing."); telemetry.events_dropped_encoder().increment(1); @@ -225,7 +230,7 @@ where pending_flush = false; - encoder.flush(context.dispatcher()).await?; + flush_encoder(&mut encoder, context.dispatcher(), &telemetry, &mut pending_events).await?; debug!("All pending payloads flushed."); } @@ -233,7 +238,28 @@ where } // Do a final flush since we may have had a pending payloads before breaking out of the loop. - encoder.flush(context.dispatcher()).await?; + flush_encoder(&mut encoder, context.dispatcher(), &telemetry, &mut pending_events).await?; + + Ok(()) +} + +/// Flushes the encoder, recording flush telemetry for the events processed since the last flush. +/// +/// Flushes with no pending events are not recorded, so the last-flush telemetry always describes a flush that sent +/// something. +async fn flush_encoder( + encoder: &mut E, dispatcher: &PayloadsDispatcher, telemetry: &ComponentTelemetry, pending_events: &mut u64, +) -> Result<(), GenericError> +where + E: IncrementalEncoder, +{ + let flush_start = Instant::now(); + encoder.flush(dispatcher).await?; + + if *pending_events > 0 { + telemetry.record_flush(*pending_events, flush_start.elapsed()); + *pending_events = 0; + } Ok(()) } diff --git a/lib/saluki-components/src/encoders/buffered_incremental/telemetry.rs b/lib/saluki-components/src/encoders/buffered_incremental/telemetry.rs index 00494bdcb34..72070dbb5ea 100644 --- a/lib/saluki-components/src/encoders/buffered_incremental/telemetry.rs +++ b/lib/saluki-components/src/encoders/buffered_incremental/telemetry.rs @@ -1,10 +1,15 @@ -use metrics::Counter; +use std::time::Duration; + +use metrics::{Counter, Gauge}; use saluki_metrics::MetricsBuilder; /// Incremental encoder-specific telemetry. #[derive(Clone)] pub struct ComponentTelemetry { events_dropped_encoder: Counter, + events_flushed: Counter, + last_flush_events: Gauge, + last_flush_duration: Gauge, } impl ComponentTelemetry { @@ -15,6 +20,9 @@ impl ComponentTelemetry { "component_events_dropped_total", ["intentional:false", "drop_reason:encoder_failure"], ), + events_flushed: builder.register_counter("encoder_flushed_events_total"), + last_flush_events: builder.register_gauge("encoder_last_flush_events"), + last_flush_duration: builder.register_gauge("encoder_last_flush_duration_nanoseconds"), } } @@ -22,4 +30,32 @@ impl ComponentTelemetry { pub fn events_dropped_encoder(&self) -> &Counter { &self.events_dropped_encoder } + + /// Records a completed flush of the given number of events, which took the given duration. + pub fn record_flush(&self, events: u64, duration: Duration) { + self.events_flushed.increment(events); + self.last_flush_events.set(events as f64); + self.last_flush_duration.set(duration.as_nanos() as f64); + } +} + +#[cfg(test)] +mod tests { + use saluki_metrics::test::TestRecorder; + + use super::*; + + #[test] + fn record_flush_tracks_totals_and_last_flush() { + let recorder = TestRecorder::default(); + let _recorder_guard = metrics::set_default_local_recorder(&recorder); + let telemetry = ComponentTelemetry::from_builder(&MetricsBuilder::default()); + + telemetry.record_flush(3, Duration::from_micros(5)); + telemetry.record_flush(2, Duration::from_micros(7)); + + assert_eq!(recorder.counter("encoder_flushed_events_total"), Some(5)); + assert_eq!(recorder.gauge("encoder_last_flush_events"), Some(2.0)); + assert_eq!(recorder.gauge("encoder_last_flush_duration_nanoseconds"), Some(7_000.0)); + } } diff --git a/lib/saluki-components/src/transforms/aggregate/mod.rs b/lib/saluki-components/src/transforms/aggregate/mod.rs index 061b8c52539..0a864512399 100644 --- a/lib/saluki-components/src/transforms/aggregate/mod.rs +++ b/lib/saluki-components/src/transforms/aggregate/mod.rs @@ -1,4 +1,9 @@ -use std::{future::pending, num::NonZeroU64, sync::Mutex, time::Duration}; +use std::{ + future::pending, + num::NonZeroU64, + sync::Mutex, + time::{Duration, Instant}, +}; use async_trait::async_trait; use ddsketch::DDSketch; @@ -26,7 +31,7 @@ use tokio::{ use tracing::{debug, error, info, trace, warn}; mod telemetry; -use self::telemetry::Telemetry; +use self::telemetry::{FlushCounts, Telemetry}; mod config; pub use self::config::HistogramConfiguration; @@ -476,6 +481,8 @@ impl Transform for Aggregate { let mut dispatcher = context.dispatcher().buffered().expect("default output should always exist"); let flush = async { + let flush_start = Instant::now(); + if let Err(e) = self.state.flush(get_unix_timestamp(), should_flush_open_windows, &mut dispatcher).await { error!(error = %e, "Failed to flush aggregation state."); } @@ -491,6 +498,8 @@ impl Transform for Aggregate { Ok(aggregated_events) => debug!(aggregated_events, "Dispatched events."), Err(e) => error!(error = %e, "Failed to flush aggregated events."), } + + self.telemetry.record_last_flush_duration(flush_start.elapsed()); }; tokio::pin!(flush); @@ -736,6 +745,8 @@ impl AggregationState { // Iterate over each context we're tracking, and flush any values that are in buckets which are now closed. debug!(timestamp = current_time, "Flushing buckets."); + let mut flush_counts = FlushCounts::default(); + for (context, am) in self.contexts.iter_mut() { // Figure out if we should remove this metric or not if it has no values in open buckets. // @@ -785,7 +796,8 @@ impl AggregationState { // This means we'll always remove all-closed/empty non-counter metrics, and we _may_ remove all-closed/empty // counters. if let Some(closed_bucket_values) = am.values.split_at_timestamp(split_timestamp) { - self.telemetry.increment_flushed(&closed_bucket_values); + self.telemetry + .increment_flushed(&closed_bucket_values, &mut flush_counts); // We got some closed bucket values, so flush those out. transform_and_push_metric( @@ -821,6 +833,7 @@ impl AggregationState { } self.last_flush = current_time; + self.telemetry.record_last_flush_counts(&flush_counts); Ok(()) } @@ -1230,6 +1243,48 @@ mod tests { } } + #[tokio::test] + async fn flush_records_last_flush_counts() { + let recorder = TestRecorder::default(); + let _local = metrics::set_default_local_recorder(&recorder); + + let builder = MetricsBuilder::default(); + let mut state = AggregationState::new( + BUCKET_WIDTH_SECS, + 10, + COUNTER_EXPIRE, + HistogramConfiguration::default(), + Telemetry::new(&builder), + ); + + assert!(state.insert(insert_ts(1), Metric::counter("metric1", 1.0))); + assert!(state.insert(insert_ts(1), Metric::gauge("metric2", 2.0))); + assert!(state.insert(insert_ts(1), Metric::distribution("metric3", 3.0))); + + let _ = get_flushed_metrics(flush_ts(1), &mut state).await; + + assert_eq!( + recorder.gauge(("aggregate_last_flush_count", &[("data_type", "series")])), + Some(2.0) + ); + assert_eq!( + recorder.gauge(("aggregate_last_flush_count", &[("data_type", "sketches")])), + Some(1.0) + ); + + // The next flush only has the idle counter's zero value, so the gauges reflect just that flush. + let _ = get_flushed_metrics(flush_ts(2), &mut state).await; + + assert_eq!( + recorder.gauge(("aggregate_last_flush_count", &[("data_type", "series")])), + Some(1.0) + ); + assert_eq!( + recorder.gauge(("aggregate_last_flush_count", &[("data_type", "sketches")])), + Some(0.0) + ); + } + #[test] fn snapshot_contexts_preserves_full_context_shape_and_unit() { let mut state = AggregationState::new( diff --git a/lib/saluki-components/src/transforms/aggregate/telemetry.rs b/lib/saluki-components/src/transforms/aggregate/telemetry.rs index f45fa65bef9..e27785384c2 100644 --- a/lib/saluki-components/src/transforms/aggregate/telemetry.rs +++ b/lib/saluki-components/src/transforms/aggregate/telemetry.rs @@ -1,3 +1,5 @@ +use std::time::Duration; + use metrics::{Counter, Gauge}; use saluki_core::data_model::event::metric::{context::Context, MetricValues}; use saluki_metrics::MetricsBuilder; @@ -48,6 +50,13 @@ impl MetricTypedGauge { } } +/// Number of series and sketches flushed in a single flush. +#[derive(Default)] +pub struct FlushCounts { + series: u64, + sketches: u64, +} + #[derive(Clone)] pub struct Telemetry { active_contexts: Gauge, @@ -57,6 +66,9 @@ pub struct Telemetry { flushes: Counter, series_flushed: Counter, sketches_flushed: Counter, + last_flush_series: Gauge, + last_flush_sketches: Gauge, + last_flush_duration: Gauge, } impl Telemetry { @@ -69,6 +81,9 @@ impl Telemetry { flushes: builder.register_counter("aggregate_flushes_total"), series_flushed: builder.register_counter_with_tags("aggregate_flushed_total", ["data_type:series"]), sketches_flushed: builder.register_counter_with_tags("aggregate_flushed_total", ["data_type:sketches"]), + last_flush_series: builder.register_gauge_with_tags("aggregate_last_flush_count", ["data_type:series"]), + last_flush_sketches: builder.register_gauge_with_tags("aggregate_last_flush_count", ["data_type:sketches"]), + last_flush_duration: builder.register_gauge("aggregate_last_flush_duration_nanoseconds"), } } @@ -82,6 +97,9 @@ impl Telemetry { flushes: Counter::noop(), series_flushed: Counter::noop(), sketches_flushed: Counter::noop(), + last_flush_series: Gauge::noop(), + last_flush_sketches: Gauge::noop(), + last_flush_duration: Gauge::noop(), } } @@ -109,11 +127,22 @@ impl Telemetry { self.flushes.increment(1); } - pub fn increment_flushed(&self, values: &MetricValues) { + pub fn increment_flushed(&self, values: &MetricValues, counts: &mut FlushCounts) { if values.is_serie() { self.series_flushed.increment(1); + counts.series += 1; } else if values.is_sketch() { self.sketches_flushed.increment(1); + counts.sketches += 1; } } + + pub fn record_last_flush_counts(&self, counts: &FlushCounts) { + self.last_flush_series.set(counts.series as f64); + self.last_flush_sketches.set(counts.sketches as f64); + } + + pub fn record_last_flush_duration(&self, duration: Duration) { + self.last_flush_duration.set(duration.as_nanos() as f64); + } } diff --git a/releasenotes/notes/rar-aggregator-flush-telemetry-943c788687171dd3.yaml b/releasenotes/notes/rar-aggregator-flush-telemetry-943c788687171dd3.yaml new file mode 100644 index 00000000000..94b09f2d308 --- /dev/null +++ b/releasenotes/notes/rar-aggregator-flush-telemetry-943c788687171dd3.yaml @@ -0,0 +1,6 @@ +enhancements: + - | + Agent Data Plane now records how many items its last aggregation and event/service check encoder + flushes handled and how long they took, and reports them through the Datadog Agent telemetry + endpoint as ``aggregator.flush_count`` and ``aggregator.flush_time``. It also reports + ``aggregator.flush`` for events and service checks.