Skip to content
Open
Show file tree
Hide file tree
Changes from 19 commits
Commits
Show all changes
23 commits
Select commit Hold shift + click to select a range
6fdd0b7
Add sample group tag to discarded events
ArunPiduguDD Sep 16, 2026
28143d7
Preserve dropped event instrumentation in sample stack
ArunPiduguDD Sep 17, 2026
473774f
Restore sample discarded event name
ArunPiduguDD Sep 21, 2026
beb7e27
Address sample metric review feedback
ArunPiduguDD Sep 22, 2026
5e02a20
Merge branch 'arun.pidugu/throttle-discarded-events-tag-group-by' int…
ArunPiduguDD Sep 29, 2026
1453ba2
Avoid untagged sample group allocation
ArunPiduguDD Sep 29, 2026
08dcd16
Merge branch 'arun.pidugu/throttle-discarded-events-tag-group-by' int…
ArunPiduguDD Sep 29, 2026
9414ee5
Merge branch 'arun.pidugu/throttle-discarded-events-tag-group-by' int…
ArunPiduguDD Sep 30, 2026
9562b26
Merge branch 'arun.pidugu/throttle-discarded-events-tag-group-by' int…
ArunPiduguDD Sep 30, 2026
a1a7823
Use tagged dropped event helper for sample
ArunPiduguDD Sep 30, 2026
0a8d0ed
Merge shared dropped-event refactor into sample
ArunPiduguDD Sep 30, 2026
09e323e
Merge branch 'arun.pidugu/throttle-discarded-events-tag-group-by' int…
ArunPiduguDD Oct 1, 2026
2649cc6
Avoid cloning retained sample groups
ArunPiduguDD Oct 1, 2026
1dd2a6b
Merge branch 'arun.pidugu/throttle-discarded-events-tag-group-by' int…
ArunPiduguDD Oct 1, 2026
15bf60d
Use group-specific dropped event helper
ArunPiduguDD Oct 1, 2026
9da0cc0
Merge branch 'arun.pidugu/throttle-discarded-events-tag-group-by' int…
ArunPiduguDD Oct 2, 2026
f05f177
Clarify sample group tag documentation
ArunPiduguDD Oct 2, 2026
a0e52cc
Restore original sample state handling
ArunPiduguDD Oct 2, 2026
954272e
Merge dropped-event test cleanup into sample
ArunPiduguDD Oct 2, 2026
bd6e657
Pass group metric setting to sample constructors
ArunPiduguDD Oct 5, 2026
40b85ce
Merge updated throttle parent into sample branch
ArunPiduguDD Oct 5, 2026
ec45e10
Merge updated throttle parent into sample branch
ArunPiduguDD Oct 6, 2026
aaa4987
Merge throttle checker fix into sample branch
ArunPiduguDD Oct 6, 2026
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
Original file line number Diff line number Diff line change
@@ -0,0 +1,3 @@
Add `internal_metrics.include_group_tag` to the `sample` transform. When enabled, `component_discarded_events_total` includes a `group` tag showing which `group_by` group an event belonged to when it was discarded. This option is disabled by default because configurations with many group values can produce a large number of unique metric tags.

authors: arun.pidugu
84 changes: 80 additions & 4 deletions src/internal_events/sample.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4,13 +4,89 @@ use vector_lib::{
};

#[derive(Debug, NamedInternalEvent)]
pub struct SampleEventDiscarded;
pub struct SampleEventDiscarded {
pub group: Option<String>,
pub include_group_tag: bool,
}

impl InternalEvent for SampleEventDiscarded {
fn emit(self) {
emit!(ComponentEventsDropped::<INTENTIONAL> {
let group = self
.include_group_tag
.then(|| self.group.unwrap_or_else(|| "None".to_string()));
ComponentEventsDropped::<INTENTIONAL> {
count: 1,
reason: "Sample discarded."
})
reason: "Sample discarded.",
}
.emit_with_group(group);
}
}

#[cfg(test)]
mod tests {
use serial_test::serial;
use vector_lib::{event::MetricValue, internal_event::InternalEvent, metrics::Controller};

use super::SampleEventDiscarded;

fn discarded_events_counter(tags: &[(&str, &str)]) -> Option<f64> {
Controller::get()
.expect("metrics controller initialized")
.capture_metrics()
.into_iter()
.find(|metric| {
metric.name() == "component_discarded_events_total"
&& tags.iter().all(|(key, value)| {
metric
.tags()
.is_some_and(|tags| tags.get(key) == Some(*value))
})
})
.map(|metric| match metric.value() {
MetricValue::Counter { value } => *value,
other => panic!("expected counter, got {other:?}"),
})
}

#[test]
#[serial]
fn emits_component_discarded_events_with_group_tag() {
vector_lib::metrics::init_test();
for group in ["group-a", "group-b"] {
SampleEventDiscarded {
group: Some(group.to_string()),
include_group_tag: true,
}
.emit();
}
SampleEventDiscarded {
group: None,
include_group_tag: true,
}
.emit();

for group in ["group-a", "group-b", "None"] {
assert_eq!(
discarded_events_counter(&[("intentional", "true"), ("group", group)]),
Some(1.0)
);
}
}

#[test]
#[serial]
fn emits_component_discarded_events_without_group_tag_by_default() {
vector_lib::metrics::init_test();
SampleEventDiscarded {
group: None,
include_group_tag: false,
}
.emit();

assert_eq!(
discarded_events_counter(&[("intentional", "true")]),
Some(1.0)
);
assert_eq!(discarded_events_counter(&[("group", "group-a")]), None);
}
}
47 changes: 46 additions & 1 deletion src/transforms/sample/config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -48,6 +48,23 @@ pub enum SampleError {
InvalidKeyFieldDynamicCombination,
}

/// Configuration of internal metrics for the Sample transform.
#[configurable_component]
#[derive(Clone, Debug, Default)]
#[serde(deny_unknown_fields)]
pub struct SampleInternalMetricsConfig {
/// Whether or not to include the `group` tag on the `component_discarded_events_total`
/// internal metric.
///
/// When enabled, adds a `group` tag containing the rendered `group_by` value.
/// Missing or unrenderable values use `None`.
///
/// Note that this defaults to false because the `group` tag has potentially unbounded
/// cardinality. Only set this to true if you know that the number of unique groups is bounded.
#[serde(default)]
pub include_group_tag: bool,
}

/// Configuration for the `sample` transform.
#[configurable_component(transform(
"sample",
Expand Down Expand Up @@ -125,6 +142,10 @@ pub struct SampleConfig {

/// A logical condition used to exclude events from sampling.
pub exclude: Option<AnyCondition>,

/// Configuration of internal metrics for the Sample transform.
#[serde(default)]
pub internal_metrics: SampleInternalMetricsConfig,
}

impl SampleConfig {
Expand Down Expand Up @@ -173,6 +194,7 @@ impl GenerateConfig for SampleConfig {
group_by: None,
exclude: None::<AnyCondition>,
sample_rate_key: default_sample_rate_key(),
internal_metrics: Default::default(),
})
.unwrap()
}
Expand Down Expand Up @@ -212,7 +234,9 @@ impl TransformConfig for SampleConfig {
)
};

Ok(Transform::function(sample))
Ok(Transform::function(sample.with_include_group_tag(
self.internal_metrics.include_group_tag,
)))
}

fn input(&self) -> Input {
Expand Down Expand Up @@ -279,6 +303,22 @@ mod tests {
crate::test_util::test_generate_config::<SampleConfig>();
}

#[test]
fn internal_metrics_include_group_tag_defaults_to_false() {
let config =
serde_yaml::from_str::<SampleConfig>("ratio: 0.5\ninternal_metrics: {}\n").unwrap();
assert!(!config.internal_metrics.include_group_tag);
}

#[test]
fn internal_metrics_include_group_tag_can_be_enabled() {
let config = serde_yaml::from_str::<SampleConfig>(
"ratio: 0.5\ninternal_metrics:\n include_group_tag: true\n",
)
.unwrap();
assert!(config.internal_metrics.include_group_tag);
}

#[test]
fn rejects_dynamic_ratio_only_configuration() {
let config = SampleConfig {
Expand All @@ -290,6 +330,7 @@ mod tests {
sample_rate_key: super::default_sample_rate_key(),
group_by: None,
exclude: None,
internal_metrics: Default::default(),
};

let err = config.sample_rate().unwrap_err();
Expand All @@ -307,6 +348,7 @@ mod tests {
sample_rate_key: super::default_sample_rate_key(),
group_by: None,
exclude: None,
internal_metrics: Default::default(),
};

let err = config.sample_rate().unwrap_err();
Expand All @@ -324,6 +366,7 @@ mod tests {
sample_rate_key: super::default_sample_rate_key(),
group_by: None,
exclude: None,
internal_metrics: Default::default(),
};

assert!(config.validate_structure().is_ok());
Expand All @@ -340,6 +383,7 @@ mod tests {
sample_rate_key: super::default_sample_rate_key(),
group_by: None,
exclude: None,
internal_metrics: Default::default(),
};

let err = config.sample_rate().unwrap_err();
Expand All @@ -357,6 +401,7 @@ mod tests {
sample_rate_key: super::default_sample_rate_key(),
group_by: None,
exclude: None,
internal_metrics: Default::default(),
};

let err = config.sample_rate().unwrap_err();
Expand Down
1 change: 1 addition & 0 deletions src/transforms/sample/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,7 @@ async fn emits_internal_events() {
group_by: None,
exclude: None,
sample_rate_key: default_sample_rate_key(),
internal_metrics: Default::default(),
};
let (tx, rx) = mpsc::channel(1);
let (topology, mut out) = create_topology(ReceiverStream::new(rx), config).await;
Expand Down
16 changes: 15 additions & 1 deletion src/transforms/sample/transform.rs
Original file line number Diff line number Diff line change
Expand Up @@ -144,6 +144,7 @@ pub struct Sample {
dynamic_event_counters: HashMap<Option<String>, u64>,
exclude: Option<Condition>,
sample_rate_key: OptionalValuePath,
include_group_tag: bool,
}

impl Sample {
Expand Down Expand Up @@ -201,9 +202,15 @@ impl Sample {
dynamic_event_counters: HashMap::default(),
exclude,
sample_rate_key,
include_group_tag: false,
}
}

pub const fn with_include_group_tag(mut self, include_group_tag: bool) -> Self {
self.include_group_tag = include_group_tag;
self
}

#[cfg(test)]
pub fn ratio(&self) -> f64 {
match &self.static_mode {
Expand Down Expand Up @@ -339,6 +346,10 @@ impl FunctionTransform for Sample {
};

let group_by_key = self.group_by_key(&event);
let discarded_group = self
.include_group_tag
.then(|| group_by_key.clone())
.flatten();
let value = self.static_key_value(&event);

let event_sample_mode = self.event_sample_mode(&event);
Expand Down Expand Up @@ -375,7 +386,10 @@ impl FunctionTransform for Sample {
}
output.push(event);
} else {
emit!(SampleEventDiscarded);
emit!(SampleEventDiscarded {
group: discarded_group,
include_group_tag: self.include_group_tag,
});
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,24 @@ generated: components: transforms: sample: configuration: {
syntax: "template"
}
}
internal_metrics: {
description: "Configuration of internal metrics for the Sample transform."
required: false
type: object: options: include_group_tag: {
description: """
Whether or not to include the `group` tag on the `component_discarded_events_total`
internal metric.

When enabled, adds a `group` tag containing the rendered `group_by` value.
Missing or unrenderable values use `None`.

Note that this defaults to false because the `group` tag has potentially unbounded
cardinality. Only set this to true if you know that the number of unique groups is bounded.
"""
required: false
type: bool: default: false
}
}
key_field: {
description: """
The name of the field whose value is hashed to determine if the event should be
Expand Down
Loading