Skip to main content

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}