vector/transforms/sample/
config.rs1use snafu::Snafu;
2use vector_lib::{
3 config::LegacyKey,
4 configurable::configurable_component,
5 lookup::{lookup_v2::OptionalValuePath, owned_value_path},
6};
7use vrl::value::Kind;
8
9use super::transform::{DynamicSampleFields, Sample, SampleMode};
10use crate::{
11 conditions::AnyCondition,
12 config::{
13 DataType, GenerateConfig, Input, OutputId, TransformConfig, TransformContext,
14 TransformOutput,
15 },
16 schema,
17 template::Template,
18 transforms::Transform,
19};
20
21#[derive(Debug, Snafu)]
22pub enum SampleError {
23 #[snafu(display(
25 "Only positive, non-zero numbers are allowed values for `ratio`, value: {ratio}"
26 ))]
27 InvalidRatio { ratio: f64 },
28
29 #[snafu(display("Only non-zero numbers are allowed values for `rate`"))]
30 InvalidRate,
31
32 #[snafu(display("Only one value can be provided for either 'rate' or 'ratio', but not both"))]
33 InvalidStaticConfiguration,
34
35 #[snafu(display(
36 "Only one value can be provided for either 'ratio_field' or 'rate_field', but not both"
37 ))]
38 InvalidDynamicConfiguration,
39
40 #[snafu(display(
41 "Exactly one value must be provided for either 'rate' or 'ratio' to configure static sampling"
42 ))]
43 MissingStaticConfiguration,
44
45 #[snafu(display(
46 "'key_field' cannot be combined with 'ratio_field' or 'rate_field' because dynamic values can vary per event and break key-based coherence"
47 ))]
48 InvalidKeyFieldDynamicCombination,
49}
50
51#[configurable_component(transform(
53 "sample",
54 "Sample events from an event stream based on supplied criteria and at a configurable rate."
55))]
56#[derive(Clone, Debug)]
57#[serde(deny_unknown_fields)]
58pub struct SampleConfig {
59 #[configurable(metadata(docs::examples = 1500))]
65 pub rate: Option<u64>,
66
67 #[configurable(metadata(docs::examples = 0.13))]
74 #[configurable(validation(range(min = 0.0, max = 1.0)))]
75 pub ratio: Option<f64>,
76
77 #[configurable(metadata(docs::examples = "sample_rate"))]
84 pub ratio_field: Option<String>,
85
86 #[configurable(metadata(docs::examples = "sample_rate_n"))]
93 pub rate_field: Option<String>,
94
95 #[configurable(metadata(docs::examples = "message"))]
109 pub key_field: Option<String>,
110
111 #[configurable(metadata(docs::examples = "sample_rate"))]
113 #[serde(default = "default_sample_rate_key")]
114 pub sample_rate_key: OptionalValuePath,
115
116 #[configurable(metadata(
124 docs::examples = "{{ service }}",
125 docs::examples = "{{ hostname }}-{{ service }}"
126 ))]
127 pub group_by: Option<Template>,
128
129 pub exclude: Option<AnyCondition>,
131}
132
133impl SampleConfig {
134 fn sample_rate(&self) -> Result<SampleMode, SampleError> {
135 if self.ratio_field.is_some() && self.rate_field.is_some() {
136 return Err(SampleError::InvalidDynamicConfiguration);
137 }
138
139 if self.key_field.is_some() && (self.ratio_field.is_some() || self.rate_field.is_some()) {
140 return Err(SampleError::InvalidKeyFieldDynamicCombination);
141 }
142
143 if self.rate.is_some() && self.ratio.is_some() {
144 return Err(SampleError::InvalidStaticConfiguration);
145 }
146
147 match (self.rate, self.ratio) {
148 (None, Some(ratio)) => {
149 if ratio <= 0.0 {
150 Err(SampleError::InvalidRatio { ratio })
151 } else {
152 Ok(SampleMode::new_ratio(ratio))
153 }
154 }
155 (Some(rate), None) => {
156 if rate == 0 {
157 Err(SampleError::InvalidRate)
158 } else {
159 Ok(SampleMode::new_rate(rate))
160 }
161 }
162 (None, None) => Err(SampleError::MissingStaticConfiguration),
163 _ => Err(SampleError::InvalidStaticConfiguration),
164 }
165 }
166}
167
168impl GenerateConfig for SampleConfig {
169 fn generate_config() -> toml::Value {
170 toml::Value::try_from(Self {
171 rate: None,
172 ratio: Some(0.1),
173 ratio_field: None,
174 rate_field: None,
175 key_field: None,
176 group_by: None,
177 exclude: None::<AnyCondition>,
178 sample_rate_key: default_sample_rate_key(),
179 })
180 .unwrap()
181 }
182}
183
184#[async_trait::async_trait]
185#[typetag::serde(name = "sample")]
186impl TransformConfig for SampleConfig {
187 async fn build(&self, context: &TransformContext) -> crate::Result<Transform> {
188 let sample_mode = self.sample_rate()?;
189 let exclude = self
190 .exclude
191 .as_ref()
192 .map(|condition| condition.build(&context.enrichment_tables, &context.metrics_storage))
193 .transpose()?;
194
195 let sample = if self.ratio_field.is_some() || self.rate_field.is_some() {
196 Sample::new_with_dynamic(
197 Self::NAME.to_string(),
198 sample_mode,
199 DynamicSampleFields {
200 ratio_field: self.ratio_field.clone(),
201 rate_field: self.rate_field.clone(),
202 },
203 self.group_by.clone(),
204 exclude,
205 self.sample_rate_key.clone(),
206 )
207 } else {
208 Sample::new(
209 Self::NAME.to_string(),
210 sample_mode,
211 self.key_field.clone(),
212 self.group_by.clone(),
213 exclude,
214 self.sample_rate_key.clone(),
215 )
216 };
217
218 Ok(Transform::function(sample))
219 }
220
221 fn input(&self) -> Input {
222 Input::new(DataType::Log | DataType::Trace)
223 }
224
225 fn validate(&self, _: &TransformContext) -> Result<(), Vec<String>> {
226 self.sample_rate()
227 .map(|_| ())
228 .map_err(|e| vec![e.to_string()])
229 }
230
231 fn validate_env(&self, context: &TransformContext) -> Result<(), Vec<String>> {
232 if let Some(Err(e)) = self
233 .exclude
234 .as_ref()
235 .map(|c| c.validate(&context.enrichment_tables, &context.metrics_storage))
236 {
237 Err(vec![format!("exclude: {e}")])
238 } else {
239 Ok(())
240 }
241 }
242
243 fn outputs(
244 &self,
245 _: &TransformContext,
246 input_definitions: &[(OutputId, schema::Definition)],
247 ) -> Vec<TransformOutput> {
248 vec![TransformOutput::new(
249 DataType::Log | DataType::Trace,
250 input_definitions
251 .iter()
252 .map(|(output, definition)| {
253 (
254 output.clone(),
255 definition.clone().with_source_metadata(
256 SampleConfig::NAME,
257 Some(LegacyKey::Overwrite(owned_value_path!("sample_rate"))),
258 &owned_value_path!("sample_rate"),
259 Kind::bytes(),
260 None,
261 ),
262 )
263 })
264 .collect(),
265 )]
266 }
267}
268
269pub fn default_sample_rate_key() -> OptionalValuePath {
270 OptionalValuePath::from(owned_value_path!("sample_rate"))
271}
272
273#[cfg(test)]
274mod tests {
275 use crate::{
276 config::TransformConfig,
277 transforms::sample::config::{SampleConfig, SampleError},
278 };
279
280 #[test]
281 fn generate_config() {
282 crate::test_util::test_generate_config::<SampleConfig>();
283 }
284
285 #[test]
286 fn rejects_dynamic_ratio_only_configuration() {
287 let config = SampleConfig {
288 rate: None,
289 ratio: None,
290 ratio_field: Some("sample_rate".to_string()),
291 rate_field: None,
292 key_field: None,
293 sample_rate_key: super::default_sample_rate_key(),
294 group_by: None,
295 exclude: None,
296 };
297
298 let err = config.sample_rate().unwrap_err();
299 assert!(matches!(err, SampleError::MissingStaticConfiguration));
300 }
301
302 #[test]
303 fn rejects_dynamic_rate_only_configuration() {
304 let config = SampleConfig {
305 rate: None,
306 ratio: None,
307 ratio_field: None,
308 rate_field: Some("sample_rate_n".to_string()),
309 key_field: None,
310 sample_rate_key: super::default_sample_rate_key(),
311 group_by: None,
312 exclude: None,
313 };
314
315 let err = config.sample_rate().unwrap_err();
316 assert!(matches!(err, SampleError::MissingStaticConfiguration));
317 }
318
319 #[test]
320 fn validates_static_with_dynamic_configuration() {
321 let config = SampleConfig {
322 rate: Some(10),
323 ratio: None,
324 ratio_field: None,
325 rate_field: Some("sample_rate_n".to_string()),
326 key_field: None,
327 sample_rate_key: super::default_sample_rate_key(),
328 group_by: None,
329 exclude: None,
330 };
331
332 assert!(
333 config
334 .validate(&crate::config::TransformContext::default())
335 .is_ok()
336 );
337 }
338
339 #[test]
340 fn rejects_both_dynamic_fields_configuration() {
341 let config = SampleConfig {
342 rate: Some(10),
343 ratio: None,
344 ratio_field: Some("sample_rate".to_string()),
345 rate_field: Some("sample_rate_n".to_string()),
346 key_field: None,
347 sample_rate_key: super::default_sample_rate_key(),
348 group_by: None,
349 exclude: None,
350 };
351
352 let err = config.sample_rate().unwrap_err();
353 assert!(matches!(err, SampleError::InvalidDynamicConfiguration));
354 }
355
356 #[test]
357 fn rejects_key_field_with_dynamic_configuration() {
358 let config = SampleConfig {
359 rate: Some(10),
360 ratio: None,
361 ratio_field: Some("sample_ratio".to_string()),
362 rate_field: None,
363 key_field: Some("trace_id".to_string()),
364 sample_rate_key: super::default_sample_rate_key(),
365 group_by: None,
366 exclude: None,
367 };
368
369 let err = config.sample_rate().unwrap_err();
370 assert!(matches!(
371 err,
372 SampleError::InvalidKeyFieldDynamicCombination
373 ));
374 }
375}