Skip to content
Merged
Changes from all commits
Commits
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
21 changes: 14 additions & 7 deletions lib/saluki-components/src/encoders/datadog/metrics/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@ use protobuf::{rt::WireType, CodedOutputStream};
use resource_accounting::{MemoryBounds, MemoryBoundsBuilder};
use saluki_common::{
buf::{ChunkedBytesBuffer, FrozenChunkedBytesBuffer},
iter::ReusableDeduplicator,
task::HandleExt as _,
};
use saluki_config::GenericConfiguration;
Expand Down Expand Up @@ -1252,9 +1253,10 @@ async fn flush_payload(
// Encodes a batch of metrics to V3 columnar format.
fn encode_v3_metrics_batch(metrics: &[Metric], additional_tags: &SharedTagSet) -> Result<Vec<u8>, GenericError> {
let mut writer = v3::V3Writer::new();
let mut tags_deduplicator = ReusableDeduplicator::new();

for metric in metrics {
write_metric_to_v3(&mut writer, metric, additional_tags);
write_metric_to_v3(&mut writer, metric, additional_tags, &mut tags_deduplicator);
Comment on lines 1255 to +1259
}

let mut output = Vec::new();
Expand All @@ -1266,7 +1268,10 @@ fn encode_v3_metrics_batch(metrics: &[Metric], additional_tags: &SharedTagSet) -
}

/// Writes a single metric to the V3 writer.
fn write_metric_to_v3(writer: &mut v3::V3Writer, metric: &Metric, additional_tags: &SharedTagSet) {
fn write_metric_to_v3(
writer: &mut v3::V3Writer, metric: &Metric, additional_tags: &SharedTagSet,
tags_deduplicator: &mut ReusableDeduplicator<Tag>,
) {
let metric_type = match metric.values() {
MetricValues::Counter(..) => v3::V3MetricType::Count,
MetricValues::Rate(..) => v3::V3MetricType::Rate,
Expand All @@ -1278,12 +1283,14 @@ fn write_metric_to_v3(writer: &mut v3::V3Writer, metric: &Metric, additional_tag
let mut builder = writer.write(metric_type, metric.context().name());

// Tags - chain instrumented + additional + origin tags
let all_tags = metric
let chained_tags = metric
.context()
.tags()
.into_iter()
.chain(additional_tags)
.chain(metric.context().origin_tags())
.chain(metric.context().origin_tags());
let all_tags = tags_deduplicator
.deduplicated(chained_tags)
.filter(|t| is_sketch || !is_v3_series_resource_tag(t) && !is_v3_series_device_tag(t))
.map(|t| t.as_str());
builder.set_tags(all_tags);
Expand All @@ -1295,13 +1302,13 @@ fn write_metric_to_v3(writer: &mut v3::V3Writer, metric: &Metric, additional_tag
}
if !is_sketch {
let mut device_resource = None;
for tag in metric
let chained_tags = metric
.context()
.origin_tags()
.into_iter()
.chain(metric.context().tags())
.chain(additional_tags)
{
.chain(additional_tags);
for tag in tags_deduplicator.deduplicated(chained_tags) {
if is_v3_series_device_tag(tag) {
device_resource = tag.value().filter(|device| !device.is_empty());
} else if is_v3_series_resource_tag(tag) {
Expand Down
Loading