Skip to content

Commit 05c30aa

Browse files
paullegranddcbryantbiggs
authored andcommitted
refactor(sdk): change SpanProcessor::on_end to take &mut FinishedSpan
Refactors the SpanProcessor API to avoid unnecessary cloning of span data: 1. Introduces `ReadableSpan` trait with getters for span fields, allowing processors to inspect spans without cloning. 2. Introduces `FinishedSpan` wrapper passed to `on_end`. Processors that only read span data pay zero cost. Processors that need ownership call `consume()` — the last processor gets the data via move (zero-copy), earlier processors get a clone. 3. Changes `on_end(&self, span: SpanData)` to `on_end(&self, span: &mut FinishedSpan)`. This is a breaking change to the SpanProcessor trait. 4. Implements `ReadableSpan` for both `Span` (live spans) and `FinishedSpan`. Benchmarks show span lifecycle cost becomes constant regardless of processor count (previously scaled linearly due to cloning). Original work by paullegranddc in PR #2962. Supersedes #2962. Relates to #2940, #2726, #2939.
1 parent 9650783 commit 05c30aa

10 files changed

Lines changed: 718 additions & 86 deletions

File tree

examples/tracing-http-propagator/src/server.rs

Lines changed: 70 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -15,12 +15,17 @@ use opentelemetry_sdk::{
1515
error::OTelSdkResult,
1616
logs::{LogProcessor, SdkLogRecord, SdkLoggerProvider},
1717
propagation::{BaggagePropagator, TraceContextPropagator},
18-
trace::{SdkTracerProvider, SpanProcessor},
18+
trace::{FinishedSpan, ReadableSpan, SdkTracerProvider, SpanProcessor},
1919
};
2020
use opentelemetry_semantic_conventions::trace;
2121
use opentelemetry_stdout::{LogExporter, SpanExporter};
2222
use std::time::Duration;
23-
use std::{convert::Infallible, net::SocketAddr, sync::OnceLock};
23+
use std::{
24+
collections::HashMap,
25+
convert::Infallible,
26+
net::SocketAddr,
27+
sync::{Mutex, OnceLock},
28+
};
2429
use tokio::net::TcpListener;
2530
use tracing::info;
2631
use tracing_subscriber::{layer::SubscriberExt, util::SubscriberInitExt};
@@ -84,6 +89,7 @@ async fn router(
8489
let span = tracer
8590
.span_builder("router")
8691
.with_kind(SpanKind::Server)
92+
.with_attributes([KeyValue::new("http.route", req.uri().path().to_string())])
8793
.start_with_context(tracer, &parent_cx);
8894

8995
info!(name = "router", message = "Dispatching request");
@@ -105,6 +111,66 @@ async fn router(
105111
response
106112
}
107113

114+
#[derive(Debug, Default)]
115+
/// A custom span processor that counts concurrent requests for each route (indentified by the http.route
116+
/// attribute) and adds that information to the span attributes.
117+
struct RouteConcurrencyCounterSpanProcessor(Mutex<HashMap<opentelemetry::Key, usize>>);
118+
119+
impl SpanProcessor for RouteConcurrencyCounterSpanProcessor {
120+
fn force_flush(&self) -> OTelSdkResult {
121+
Ok(())
122+
}
123+
124+
fn shutdown_with_timeout(&self, _timeout: Duration) -> crate::OTelSdkResult {
125+
Ok(())
126+
}
127+
128+
fn on_start(&self, span: &mut opentelemetry_sdk::trace::Span, _cx: &Context) {
129+
if !matches!(span.span_kind(), SpanKind::Server) {
130+
return;
131+
}
132+
let Some(route) = span
133+
.attributes()
134+
.iter()
135+
.find(|kv| kv.key.as_str() == "http.route")
136+
else {
137+
return;
138+
};
139+
let Ok(mut counts) = self.0.lock() else {
140+
return;
141+
};
142+
let count = counts.entry(route.key.clone()).or_default();
143+
*count += 1;
144+
span.set_attribute(KeyValue::new(
145+
"http.route.concurrent_requests",
146+
*count as i64,
147+
));
148+
}
149+
150+
fn on_end(&self, span: &mut FinishedSpan) {
151+
if !matches!(span.span_kind(), SpanKind::Server) {
152+
return;
153+
}
154+
let Some(route) = span
155+
.attributes()
156+
.iter()
157+
.find(|kv| kv.key.as_str() == "http.route")
158+
else {
159+
return;
160+
};
161+
let Ok(mut counts) = self.0.lock() else {
162+
return;
163+
};
164+
let Some(count) = counts.get_mut(&route.key) else {
165+
return;
166+
};
167+
*count -= 1;
168+
if *count == 0 {
169+
counts.remove(&route.key);
170+
}
171+
}
172+
}
173+
108174
/// A custom log processor that enriches LogRecords with baggage attributes.
109175
/// Baggage information is not added automatically without this processor.
110176
#[derive(Debug)]
@@ -146,7 +212,7 @@ impl SpanProcessor for EnrichWithBaggageSpanProcessor {
146212
}
147213
}
148214

149-
fn on_end(&self, _span: opentelemetry_sdk::trace::SpanData) {}
215+
fn on_end(&self, _span: &mut opentelemetry_sdk::trace::FinishedSpan) {}
150216
}
151217

152218
fn init_tracer() -> SdkTracerProvider {
@@ -162,6 +228,7 @@ fn init_tracer() -> SdkTracerProvider {
162228
// Setup tracerprovider with stdout exporter
163229
// that prints the spans to stdout.
164230
let provider = SdkTracerProvider::builder()
231+
.with_span_processor(RouteConcurrencyCounterSpanProcessor::default())
165232
.with_span_processor(EnrichWithBaggageSpanProcessor)
166233
.with_simple_exporter(SpanExporter::default())
167234
.build();

opentelemetry-sdk/Cargo.toml

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -125,6 +125,12 @@ name = "log"
125125
harness = false
126126
required-features = ["logs"]
127127

128+
[[bench]]
129+
name = "span_processor_api"
130+
harness = false
131+
required-features = ["testing"]
132+
133+
128134
[lib]
129135
bench = false
130136

opentelemetry-sdk/benches/batch_span_processor.rs

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -4,10 +4,10 @@ use opentelemetry::trace::{
44
SpanContext, SpanId, SpanKind, Status, TraceFlags, TraceId, TraceState,
55
};
66
use opentelemetry_sdk::testing::trace::NoopSpanExporter;
7-
use opentelemetry_sdk::trace::SpanData;
87
use opentelemetry_sdk::trace::{
98
BatchConfigBuilder, BatchSpanProcessor, SpanEvents, SpanLinks, SpanProcessor,
109
};
10+
use opentelemetry_sdk::trace::{FinishedSpan, SpanData};
1111
use std::sync::Arc;
1212
use tokio::runtime::Runtime;
1313

@@ -63,7 +63,8 @@ fn criterion_benchmark(c: &mut Criterion) {
6363
let spans = get_span_data();
6464
handles.push(tokio::spawn(async move {
6565
for span in spans {
66-
span_processor.on_end(span);
66+
let mut span = FinishedSpan::new(span);
67+
span_processor.on_end(&mut span);
6768
tokio::task::yield_now().await;
6869
}
6970
}));
Lines changed: 88 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,88 @@
1+
use std::time::Duration;
2+
3+
use criterion::{black_box, criterion_group, criterion_main, Criterion};
4+
use opentelemetry::{
5+
trace::{Span, Tracer, TracerProvider},
6+
Context, KeyValue,
7+
};
8+
use opentelemetry_sdk::trace as sdktrace;
9+
10+
#[cfg(not(target_os = "windows"))]
11+
use pprof::criterion::{Output, PProfProfiler};
12+
13+
/*
14+
Adding results in comments for a quick reference.
15+
Chip: Apple M1 Max
16+
Total Number of Cores: 10 (8 performance and 2 efficiency)
17+
18+
SpanProcessorApi/0_processors
19+
time: [385.15 ns 386.14 ns 387.25 ns]
20+
SpanProcessorApi/1_processors
21+
time: [385.73 ns 387.17 ns 388.85 ns]
22+
SpanProcessorApi/2_processors
23+
time: [384.84 ns 385.66 ns 386.50 ns]
24+
SpanProcessorApi/4_processors
25+
time: [386.78 ns 388.17 ns 389.58 ns]
26+
*/
27+
28+
#[derive(Debug)]
29+
struct NoopSpanProcessor;
30+
31+
impl sdktrace::SpanProcessor for NoopSpanProcessor {
32+
fn on_start(&self, _span: &mut sdktrace::Span, _parent_cx: &Context) {}
33+
fn on_end(&self, _span: &mut sdktrace::FinishedSpan) {}
34+
fn force_flush(&self) -> opentelemetry_sdk::error::OTelSdkResult {
35+
Ok(())
36+
}
37+
fn shutdown_with_timeout(&self, _timeout: Duration) -> opentelemetry_sdk::error::OTelSdkResult {
38+
Ok(())
39+
}
40+
}
41+
42+
fn create_tracer(span_processors_count: usize) -> sdktrace::SdkTracer {
43+
let mut builder = sdktrace::SdkTracerProvider::builder();
44+
for _ in 0..span_processors_count {
45+
builder = builder.with_span_processor(NoopSpanProcessor);
46+
}
47+
builder.build().tracer("tracer")
48+
}
49+
50+
fn create_span(tracer: &sdktrace::Tracer) -> sdktrace::Span {
51+
let mut span = tracer.start("foo");
52+
span.set_attribute(KeyValue::new("key1", false));
53+
span.set_attribute(KeyValue::new("key2", "hello"));
54+
span.set_attribute(KeyValue::new("key4", 123.456));
55+
span.add_event("my_event", vec![KeyValue::new("key1", "value1")]);
56+
span
57+
}
58+
59+
fn criterion_benchmark(c: &mut Criterion) {
60+
let mut group = c.benchmark_group("SpanProcessorApi");
61+
for i in [0, 1, 2, 4] {
62+
group.bench_function(format!("{}_processors", i), |b| {
63+
let tracer = create_tracer(i);
64+
b.iter(|| {
65+
black_box(create_span(&tracer));
66+
});
67+
});
68+
}
69+
}
70+
71+
#[cfg(not(target_os = "windows"))]
72+
criterion_group! {
73+
name = benches;
74+
config = Criterion::default().with_profiler(PProfProfiler::new(100, Output::Flamegraph(None)))
75+
.warm_up_time(std::time::Duration::from_secs(1))
76+
.measurement_time(std::time::Duration::from_secs(2));
77+
targets = criterion_benchmark
78+
}
79+
80+
#[cfg(target_os = "windows")]
81+
criterion_group! {
82+
name = benches;
83+
config = Criterion::default().warm_up_time(std::time::Duration::from_secs(1))
84+
.measurement_time(std::time::Duration::from_secs(2));
85+
targets = criterion_benchmark
86+
}
87+
88+
criterion_main!(benches);

opentelemetry-sdk/src/trace/mod.rs

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -39,7 +39,7 @@ pub use id_generator::{IdGenerator, RandomIdGenerator};
3939
pub use links::SpanLinks;
4040
pub use provider::{SdkTracerProvider, TracerProviderBuilder};
4141
pub use sampler::{Sampler, SamplingDecision, SamplingResult, ShouldSample};
42-
pub use span::Span;
42+
pub use span::{FinishedSpan, ReadableSpan, Span};
4343
pub use span_limit::SpanLimits;
4444
pub use span_processor::{
4545
BatchConfig, BatchConfigBuilder, BatchSpanProcessor, BatchSpanProcessorBuilder,
@@ -138,7 +138,7 @@ mod tests {
138138
}
139139
}
140140

141-
fn on_end(&self, span: SpanData) {
141+
fn on_end(&self, span: &mut FinishedSpan) {
142142
// Fixed: Context::current() no longer panics from Drop
143143
// See https://github.com/open-telemetry/opentelemetry-rust/issues/2871
144144
Context::current();
@@ -149,7 +149,7 @@ mod tests {
149149

150150
// Verify: on_start stored the baggage as an attribute
151151
assert!(
152-
span.attributes
152+
span.attributes()
153153
.iter()
154154
.any(|kv| kv.key.as_str() == "bag-key"),
155155
"Baggage should have been stored as span attribute in on_start"

opentelemetry-sdk/src/trace/provider.rs

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -510,8 +510,8 @@ mod tests {
510510
SERVICE_NAME, TELEMETRY_SDK_LANGUAGE, TELEMETRY_SDK_NAME, TELEMETRY_SDK_VERSION,
511511
};
512512
use crate::trace::provider::TracerProviderInner;
513-
use crate::trace::{Config, Span, SpanProcessor};
514-
use crate::trace::{SdkTracerProvider, SpanData};
513+
use crate::trace::SdkTracerProvider;
514+
use crate::trace::{Config, FinishedSpan, Span, SpanProcessor};
515515
use crate::Resource;
516516
use opentelemetry::trace::{Tracer, TracerProvider};
517517
use opentelemetry::{Context, Key, KeyValue, Value};
@@ -565,7 +565,7 @@ mod tests {
565565
.fetch_add(1, Ordering::SeqCst);
566566
}
567567

568-
fn on_end(&self, _span: SpanData) {
568+
fn on_end(&self, _span: &mut FinishedSpan) {
569569
// ignore
570570
}
571571

@@ -858,7 +858,7 @@ mod tests {
858858
// No operation needed for this processor
859859
}
860860

861-
fn on_end(&self, _span: SpanData) {
861+
fn on_end(&self, _span: &mut FinishedSpan) {
862862
// No operation needed for this processor
863863
}
864864

0 commit comments

Comments
 (0)