relay_metrics/bucket.rs
1use std::collections::{BTreeMap, BTreeSet};
2use std::{fmt, mem};
3
4use relay_common::time::UnixTimestamp;
5use relay_protocol::FiniteF64;
6use serde::{Deserialize, Serialize};
7use smallvec::SmallVec;
8
9use crate::protocol::{
10 CounterType, DistributionType, GaugeType, MetricName, MetricType, SetType, hash_set_value,
11};
12
13/// Type of [`Bucket::tags`].
14pub type MetricTags = BTreeMap<String, String>;
15
16/// A snapshot of values within a [`Bucket`].
17#[derive(Clone, Copy, Debug, PartialEq, Deserialize, Serialize)]
18pub struct GaugeValue {
19 /// The last value reported in the bucket.
20 ///
21 /// This aggregation is not commutative.
22 pub last: GaugeType,
23 /// The minimum value reported in the bucket.
24 pub min: GaugeType,
25 /// The maximum value reported in the bucket.
26 pub max: GaugeType,
27 /// The sum of all values reported in the bucket.
28 pub sum: GaugeType,
29 /// The number of times this bucket was updated with a new value.
30 pub count: u64,
31}
32
33impl GaugeValue {
34 /// Creates a gauge snapshot from a single value.
35 pub fn single(value: GaugeType) -> Self {
36 Self {
37 last: value,
38 min: value,
39 max: value,
40 sum: value,
41 count: 1,
42 }
43 }
44
45 /// Inserts a new value into the gauge.
46 pub fn insert(&mut self, value: GaugeType) {
47 self.last = value;
48 self.min = self.min.min(value);
49 self.max = self.max.max(value);
50 self.sum = self.sum.saturating_add(value);
51 self.count += 1;
52 }
53
54 /// Merges two gauge snapshots.
55 pub fn merge(&mut self, other: Self) {
56 self.last = other.last;
57 self.min = self.min.min(other.min);
58 self.max = self.max.max(other.max);
59 self.sum = self.sum.saturating_add(other.sum);
60 self.count += other.count;
61 }
62
63 /// Returns the average of all values reported in this bucket.
64 pub fn avg(&self) -> Option<GaugeType> {
65 self.sum / FiniteF64::new(self.count as f64)?
66 }
67}
68
69/// A distribution of values within a [`Bucket`].
70///
71/// Distributions logically store a histogram of values. Based on individual reported values,
72/// distributions allow to query the maximum, minimum, or average of the reported values, as well as
73/// statistical quantiles.
74///
75/// # Example
76///
77/// ```
78/// use relay_metrics::dist;
79///
80/// let mut dist = dist![1, 1, 1, 2];
81/// dist.push(5.into());
82/// dist.extend(std::iter::repeat(3.into()).take(7));
83/// ```
84///
85/// Logically, this distribution is equivalent to this visualization:
86///
87/// ```plain
88/// value | count
89/// 1.0 | ***
90/// 2.0 | *
91/// 3.0 | *******
92/// 4.0 |
93/// 5.0 | *
94/// ```
95///
96/// # Serialization
97///
98/// Distributions serialize as lists of floating point values. The list contains one entry for each
99/// value in the distribution, including duplicates.
100pub type DistributionValue = SmallVec<[DistributionType; 3]>;
101
102#[doc(hidden)]
103pub use smallvec::smallvec as _smallvec;
104
105/// Creates a [`DistributionValue`] containing the given arguments.
106///
107/// `dist!` allows `DistributionValue` to be defined with the same syntax as array expressions.
108///
109/// # Example
110///
111/// ```
112/// let dist = relay_metrics::dist![1, 2];
113/// ```
114#[macro_export]
115macro_rules! dist {
116 ($($x:expr),*$(,)*) => {
117 $crate::_smallvec!($($crate::DistributionType::from($x)),*) as $crate::DistributionValue
118 };
119}
120
121/// A set of unique values.
122///
123/// Values are represented as 32-bit integer hashes. Use [`BucketValue::set_from_str`] to hash a
124/// string value before adding it to a set.
125pub type SetValue = BTreeSet<SetType>;
126
127/// The [aggregated value](Bucket::value) of a metric bucket.
128#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
129#[serde(tag = "type", content = "value")]
130pub enum BucketValue {
131 /// Counts instances of an event ([`MetricType::Counter`]).
132 ///
133 /// Counters can be incremented and decremented. The default operation is to increment a counter
134 /// by `1`, although increments by larger values are equally possible.
135 ///
136 /// # Serialization
137 ///
138 /// This variant serializes to a double precision float.
139 ///
140 /// # Aggregation
141 ///
142 /// Counters aggregate by folding individual values into a single sum value per bucket. The sum
143 /// is ingested and stored directly.
144 #[serde(rename = "c")]
145 Counter(CounterType),
146
147 /// Builds a statistical distribution over values reported ([`MetricType::Distribution`]).
148 ///
149 /// Based on individual reported values, distributions allow to query the maximum, minimum, or
150 /// average of the reported values, as well as statistical quantiles. With an increasing number
151 /// of values in the distribution, its accuracy becomes approximate.
152 ///
153 /// # Serialization
154 ///
155 /// This variant serializes to a list of double precision floats, see [`DistributionValue`].
156 ///
157 /// # Aggregation
158 ///
159 /// During ingestion, all individual reported values are collected in a lossless format. In
160 /// storage, these values are compressed into data sketches that allow to query quantiles.
161 /// Separately, the count and sum of the reported values is stored, which makes distributions a
162 /// strict superset of counters.
163 #[serde(rename = "d")]
164 Distribution(DistributionValue),
165
166 /// Counts the number of unique reported values.
167 ///
168 /// Sets represent arbitrary discrete values as hashes and track their deduplicated count. With
169 /// an increasing number of unique values in the set, its accuracy becomes approximate. It is
170 /// not possible to query individual values from a set.
171 ///
172 /// # Serialization
173 ///
174 /// This variant serializes to a list of 32-bit integers.
175 ///
176 /// # Aggregation
177 ///
178 /// Set values are represented as 32-bit integer hashes of the original value. Buckets merge
179 /// by taking the union of these hashes.
180 ///
181 /// Internally, set metrics are stored in data sketches that expose an approximate cardinality.
182 #[serde(rename = "s")]
183 Set(SetValue),
184
185 /// Stores absolute snapshots of values.
186 ///
187 /// In addition to plain [counters](Self::Counter), gauges store a snapshot of the maximum,
188 /// minimum and sum of all values, as well as the last reported value. Note that the "last"
189 /// component of this aggregation is not commutative. Which value is preserved as last value is
190 /// implementation-defined.
191 ///
192 /// # Serialization
193 ///
194 /// This variant serializes to a structure with named fields, see [`GaugeValue`].
195 ///
196 /// # Aggregation
197 ///
198 /// Gauges aggregate by folding each of the components based on their semantics:
199 /// - `last` assumes the newly added value
200 /// - `min` retains the smaller value
201 /// - `max` retains the larger value
202 /// - `sum` adds the new value to the existing sum
203 /// - `count` adds the count of the newly added gauge (defaulting to `1`)
204 #[serde(rename = "g")]
205 Gauge(GaugeValue),
206}
207
208impl BucketValue {
209 /// Returns a bucket value representing a counter with the given value.
210 pub fn counter(value: CounterType) -> Self {
211 Self::Counter(value)
212 }
213
214 /// Returns a bucket value representing a distribution with a single given value.
215 pub fn distribution(value: DistributionType) -> Self {
216 Self::Distribution(dist![value])
217 }
218
219 /// Returns a bucket value representing a set with a single given hash value.
220 pub fn set(value: SetType) -> Self {
221 Self::Set(std::iter::once(value).collect())
222 }
223
224 /// Returns a bucket value representing a set with the 32-bit hash of the given string.
225 pub fn set_from_str(string: &str) -> Self {
226 Self::set(hash_set_value(string))
227 }
228
229 /// Returns a bucket value representing a set with a single given value.
230 pub fn set_from_display(display: impl fmt::Display) -> Self {
231 Self::set(hash_set_value(&display.to_string()))
232 }
233
234 /// Returns a bucket value representing a gauge with a single given value.
235 pub fn gauge(value: GaugeType) -> Self {
236 Self::Gauge(GaugeValue::single(value))
237 }
238
239 /// Returns the type of this value.
240 pub fn ty(&self) -> MetricType {
241 match self {
242 Self::Counter(_) => MetricType::Counter,
243 Self::Distribution(_) => MetricType::Distribution,
244 Self::Set(_) => MetricType::Set,
245 Self::Gauge(_) => MetricType::Gauge,
246 }
247 }
248
249 /// Returns the number of raw data points in this value.
250 pub fn len(&self) -> usize {
251 match self {
252 BucketValue::Counter(_) => 1,
253 BucketValue::Distribution(distribution) => distribution.len(),
254 BucketValue::Set(set) => set.len(),
255 BucketValue::Gauge(_) => 5,
256 }
257 }
258
259 /// Returns `true` if this bucket contains no values.
260 pub fn is_empty(&self) -> bool {
261 self.len() == 0
262 }
263
264 /// Estimates the number of bytes needed to encode the bucket value.
265 ///
266 /// Note that this does not necessarily match the exact memory footprint of the value,
267 /// because data structures have a memory overhead.
268 pub fn cost(&self) -> usize {
269 // Beside the size of [`BucketValue`], we also need to account for the cost of values
270 // allocated dynamically.
271 let allocated_cost = match self {
272 Self::Counter(_) => 0,
273 Self::Set(s) => mem::size_of::<SetType>() * s.len(),
274 Self::Gauge(_) => 0,
275 Self::Distribution(d) => d.len() * mem::size_of::<DistributionType>(),
276 };
277
278 mem::size_of::<Self>() + allocated_cost
279 }
280
281 /// Merges the given `bucket_value` into `self`.
282 ///
283 /// Returns `Ok(())` if the two bucket values can be merged. This is the case when both bucket
284 /// values are of the same variant. Otherwise, this returns `Err(other)`.
285 pub fn merge(&mut self, other: Self) -> Result<(), Self> {
286 match (self, other) {
287 (Self::Counter(slf), Self::Counter(other)) => *slf = slf.saturating_add(other),
288 (Self::Distribution(slf), Self::Distribution(other)) => slf.extend_from_slice(&other),
289 (Self::Set(slf), Self::Set(other)) => slf.extend(other),
290 (Self::Gauge(slf), Self::Gauge(other)) => slf.merge(other),
291 (_, other) => return Err(other),
292 }
293
294 Ok(())
295 }
296}
297
298/// An aggregation of metric values.
299///
300/// As opposed to single metric values, bucket aggregations can carry multiple values. See
301/// [`MetricType`] for a description on how values are aggregated in buckets. Values are aggregated
302/// by metric name, type, time window, and all tags. Particularly, this allows metrics to have the
303/// same name even if their types differ.
304///
305/// See the [crate documentation](crate) for general information on Metrics.
306///
307/// # Values
308///
309/// The contents of a bucket, especially their representation and serialization, depend on the
310/// metric type:
311///
312/// - [Counters](BucketValue::Counter) store a single value, serialized as floating point.
313/// - [Distributions](MetricType::Distribution) and [sets](MetricType::Set) store the full set of
314/// reported values.
315/// - [Gauges](BucketValue::Gauge) store a snapshot of reported values, see [`GaugeValue`].
316///
317/// # JSON Representation
318///
319/// Metrics are represented as structured data in JSON.
320/// The data type of the `value` field is determined by the metric type.
321///
322/// Buckets have a required [`width`](Self::width) field in their JSON representation.
323///
324/// ```json
325#[doc = include_str!("../tests/fixtures/buckets.json")]
326/// ```
327#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
328pub struct Bucket {
329 /// The start time of the bucket's time window, in seconds since the Unix epoch.
330 pub timestamp: UnixTimestamp,
331
332 /// The length of the time window in seconds.
333 ///
334 /// To initialize a new bucket, choose `0` as width. Once the bucket is tracked by Relay's
335 /// aggregator, the width is aligned with configuration for the namespace and the timestamp is
336 /// adjusted accordingly.
337 pub width: u64,
338
339 /// The name of the metric in MRI (metric resource identifier) format.
340 ///
341 /// MRIs have the format `<type>:<ns>/<name>@<unit>`. See [`crate::MetricResourceIdentifier`] for
342 /// information on fields and representations.
343 pub name: MetricName,
344
345 /// The type and aggregated values of this bucket.
346 ///
347 /// The metric type determines how values are aggregated and represented in JSON.
348 ///
349 /// See [`BucketValue`] for more examples and semantics.
350 #[serde(flatten)]
351 pub value: BucketValue,
352
353 /// A list of tags adding dimensions to the metric for filtering and aggregation.
354 ///
355 /// Tags allow to compute separate aggregates to filter or group metric values by any number of
356 /// dimensions. Tags consist of a unique tag key and one associated value. For tags with missing
357 /// values, an empty `""` value is assumed at query time.
358 #[serde(default, skip_serializing_if = "MetricTags::is_empty")]
359 pub tags: MetricTags,
360
361 /// Relay internal metadata for a metric bucket.
362 ///
363 /// The metadata contains meta information about the metric bucket itself,
364 /// for example how many this bucket has been aggregated in total.
365 #[serde(default, skip_serializing_if = "BucketMetadata::is_default")]
366 pub metadata: BucketMetadata,
367}
368
369impl Bucket {
370 /// Returns the value of the specified tag if it exists.
371 pub fn tag(&self, name: &str) -> Option<&str> {
372 self.tags.get(name).map(|s| s.as_str())
373 }
374
375 /// Removes the value of the specified tag.
376 ///
377 /// If the tag exists, the removed value is returned.
378 pub fn remove_tag(&mut self, name: &str) -> Option<String> {
379 self.tags.remove(name)
380 }
381}
382
383/// Relay internal metadata for a metric bucket.
384#[derive(Clone, Copy, Debug, PartialEq, Deserialize, Serialize)]
385pub struct BucketMetadata {
386 /// How many times the bucket was merged.
387 ///
388 /// Creating a new bucket is the first merge.
389 /// Merging two buckets sums the amount of merges.
390 ///
391 /// For example: Merging two un-merged buckets will yield a total
392 /// of `2` merges.
393 ///
394 /// Due to how Relay aggregates metrics and later splits them into multiple
395 /// buckets again, the amount of merges can be zero.
396 /// When splitting a bucket the total volume of the bucket may only be attributed
397 /// to one part or distributed across the resulting buckets, in either case
398 /// values of `0` are possible.
399 pub merges: u32,
400
401 /// Received timestamp of the first metric in this bucket.
402 ///
403 /// This field should be set to the time in which the first metric of a specific bucket was
404 /// received in the outermost internal Relay.
405 pub received_at: Option<UnixTimestamp>,
406
407 /// Is `true` if this metric was extracted from a sampled/indexed envelope item.
408 ///
409 /// The final dynamic sampling decision is always made in processing Relays.
410 /// If a metric was extracted from an item which is sampled (i.e. retained by dynamic sampling), this flag is `true`.
411 ///
412 /// Since these metrics from samples carry additional information, e.g. they don't
413 /// require rate limiting since the sample they've been extracted from was already
414 /// rate limited, this flag must be included in the aggregation key when aggregation buckets.
415 #[serde(skip)]
416 pub extracted_from_indexed: bool,
417}
418
419impl BucketMetadata {
420 /// Creates a fresh metadata instance.
421 ///
422 /// The new metadata is initialized with `1` merge and a given `received_at` timestamp.
423 pub fn new(received_at: UnixTimestamp) -> Self {
424 Self {
425 merges: 1,
426 received_at: Some(received_at),
427 extracted_from_indexed: false,
428 }
429 }
430
431 /// Whether the metadata does not contain more information than the default.
432 pub fn is_default(&self) -> bool {
433 &Self::default() == self
434 }
435
436 /// Merges another metadata object into the current one.
437 pub fn merge(&mut self, other: Self) {
438 self.merges = self.merges.saturating_add(other.merges);
439 self.received_at = match (self.received_at, other.received_at) {
440 (Some(received_at), None) => Some(received_at),
441 (None, Some(received_at)) => Some(received_at),
442 (left, right) => left.min(right),
443 };
444 }
445}
446
447impl Default for BucketMetadata {
448 fn default() -> Self {
449 Self {
450 merges: 1,
451 received_at: None,
452 extracted_from_indexed: false,
453 }
454 }
455}
456
457#[cfg(test)]
458mod tests {
459 use similar_asserts::assert_eq;
460
461 use super::*;
462
463 #[test]
464 fn test_distribution_value_size() {
465 // DistributionValue uses a SmallVec internally to prevent an additional allocation and
466 // indirection in cases where it needs to store only a small number of items. This is
467 // enabled by a comparably large `GaugeValue`, which stores five atoms. Ensure that the
468 // `DistributionValue`'s size does not exceed that of `GaugeValue`.
469 assert!(
470 std::mem::size_of::<DistributionValue>() <= std::mem::size_of::<GaugeValue>(),
471 "distribution value should not exceed gauge {}",
472 std::mem::size_of::<DistributionValue>()
473 );
474 }
475
476 #[test]
477 fn test_bucket_value_merge_counter() {
478 let mut value = BucketValue::Counter(42.into());
479 value.merge(BucketValue::Counter(43.into())).unwrap();
480 assert_eq!(value, BucketValue::Counter(85.into()));
481 }
482
483 #[test]
484 fn test_bucket_value_merge_distribution() {
485 let mut value = BucketValue::Distribution(dist![1, 2, 3]);
486 value.merge(BucketValue::Distribution(dist![2, 4])).unwrap();
487 assert_eq!(value, BucketValue::Distribution(dist![1, 2, 3, 2, 4]));
488 }
489
490 #[test]
491 fn test_bucket_value_merge_set() {
492 let mut value = BucketValue::Set(vec![1, 2].into_iter().collect());
493 value.merge(BucketValue::Set([2, 3].into())).unwrap();
494 assert_eq!(value, BucketValue::Set(vec![1, 2, 3].into_iter().collect()));
495 }
496
497 #[test]
498 fn test_bucket_value_merge_gauge() {
499 let mut value = BucketValue::Gauge(GaugeValue::single(42.into()));
500 value.merge(BucketValue::gauge(43.into())).unwrap();
501
502 assert_eq!(
503 value,
504 BucketValue::Gauge(GaugeValue {
505 last: 43.into(),
506 min: 42.into(),
507 max: 43.into(),
508 sum: 85.into(),
509 count: 2,
510 })
511 );
512 }
513
514 #[test]
515 fn test_parse_buckets() {
516 let json = r#"[
517 {
518 "name": "endpoint.response_time",
519 "unit": "millisecond",
520 "value": [36, 49, 57, 68],
521 "type": "d",
522 "timestamp": 1615889440,
523 "width": 10,
524 "tags": {
525 "route": "user_index"
526 },
527 "metadata": {
528 "merges": 1,
529 "received_at": 1615889440
530 }
531 }
532 ]"#;
533
534 let buckets = serde_json::from_str::<Vec<Bucket>>(json).unwrap();
535
536 insta::assert_debug_snapshot!(buckets, @r###"
537 [
538 Bucket {
539 timestamp: UnixTimestamp(1615889440),
540 width: 10,
541 name: MetricName(
542 "endpoint.response_time",
543 ),
544 value: Distribution(
545 [
546 36.0,
547 49.0,
548 57.0,
549 68.0,
550 ],
551 ),
552 tags: {
553 "route": "user_index",
554 },
555 metadata: BucketMetadata {
556 merges: 1,
557 received_at: Some(
558 UnixTimestamp(1615889440),
559 ),
560 extracted_from_indexed: false,
561 },
562 },
563 ]
564 "###);
565 }
566
567 #[test]
568 fn test_parse_bucket_defaults() {
569 let json = r#"[
570 {
571 "name": "endpoint.hits",
572 "value": 4,
573 "type": "c",
574 "timestamp": 1615889440,
575 "width": 10,
576 "metadata": {
577 "merges": 1,
578 "received_at": 1615889440
579 }
580 }
581 ]"#;
582
583 let buckets = serde_json::from_str::<Vec<Bucket>>(json).unwrap();
584
585 insta::assert_debug_snapshot!(buckets, @r###"
586 [
587 Bucket {
588 timestamp: UnixTimestamp(1615889440),
589 width: 10,
590 name: MetricName(
591 "endpoint.hits",
592 ),
593 value: Counter(
594 4.0,
595 ),
596 tags: {},
597 metadata: BucketMetadata {
598 merges: 1,
599 received_at: Some(
600 UnixTimestamp(1615889440),
601 ),
602 extracted_from_indexed: false,
603 },
604 },
605 ]
606 "###);
607 }
608
609 #[test]
610 fn test_buckets_roundtrip() {
611 let json = r#"[
612 {
613 "timestamp": 1615889440,
614 "width": 10,
615 "name": "endpoint.response_time",
616 "type": "d",
617 "value": [
618 36.0,
619 49.0,
620 57.0,
621 68.0
622 ],
623 "tags": {
624 "route": "user_index"
625 }
626 },
627 {
628 "timestamp": 1615889440,
629 "width": 10,
630 "name": "endpoint.hits",
631 "type": "c",
632 "value": 4.0,
633 "tags": {
634 "route": "user_index"
635 }
636 },
637 {
638 "timestamp": 1615889440,
639 "width": 10,
640 "name": "endpoint.parallel_requests",
641 "type": "g",
642 "value": {
643 "last": 25.0,
644 "min": 17.0,
645 "max": 42.0,
646 "sum": 2210.0,
647 "count": 85
648 }
649 },
650 {
651 "timestamp": 1615889440,
652 "width": 10,
653 "name": "endpoint.users",
654 "type": "s",
655 "value": [
656 3182887624,
657 4267882815
658 ],
659 "tags": {
660 "route": "user_index"
661 }
662 }
663]"#;
664
665 let buckets = serde_json::from_str::<Vec<Bucket>>(json).unwrap();
666 let serialized = serde_json::to_string_pretty(&buckets).unwrap();
667 assert_eq!(json, serialized);
668 }
669
670 #[test]
671 fn test_bucket_docs_roundtrip() {
672 let json = include_str!("../tests/fixtures/buckets.json")
673 .trim_end()
674 .replace("\r\n", "\n");
675 let buckets = serde_json::from_str::<Vec<Bucket>>(&json).unwrap();
676
677 let serialized = serde_json::to_string_pretty(&buckets).unwrap();
678 assert_eq!(json, serialized);
679 }
680
681 #[test]
682 fn test_bucket_metadata_merge() {
683 let mut metadata = BucketMetadata::default();
684
685 let other_metadata = BucketMetadata::default();
686 metadata.merge(other_metadata);
687 assert_eq!(
688 metadata,
689 BucketMetadata {
690 merges: 2,
691 received_at: None,
692 extracted_from_indexed: false,
693 }
694 );
695
696 let other_metadata = BucketMetadata::new(UnixTimestamp::from_secs(10));
697 metadata.merge(other_metadata);
698 assert_eq!(
699 metadata,
700 BucketMetadata {
701 merges: 3,
702 received_at: Some(UnixTimestamp::from_secs(10)),
703 extracted_from_indexed: false,
704 }
705 );
706
707 let other_metadata = BucketMetadata::new(UnixTimestamp::from_secs(20));
708 metadata.merge(other_metadata);
709 assert_eq!(
710 metadata,
711 BucketMetadata {
712 merges: 4,
713 received_at: Some(UnixTimestamp::from_secs(10)),
714 extracted_from_indexed: false,
715 }
716 );
717 }
718}