Skip to main content

SpanProcessor

Trait SpanProcessor 

pub trait SpanProcessor:
    Send
    + Sync
    + Debug {
    // Required methods
    fn on_start(&self, span: &mut Span, cx: &Context);
    fn on_end(&self, span: SpanData);
    fn force_flush(&self) -> Result<(), OTelSdkError>;
    fn shutdown_with_timeout(
        &self,
        timeout: Duration,
    ) -> Result<(), OTelSdkError>;

    // Provided methods
    fn shutdown(&self) -> Result<(), OTelSdkError> { ... }
    fn set_resource(&mut self, _resource: &Resource) { ... }
}
Available on crate features opentelemetry and trace only.
Expand description

SpanProcessor is an interface which allows hooks for span start and end method invocations. The span processors are invoked only when is_recording is true.

Required Methods§

fn on_start(&self, span: &mut Span, cx: &Context)

on_start is called when a Span is started. This method is called synchronously on the thread that started the span, therefore it should not block or throw exceptions.

fn on_end(&self, span: SpanData)

on_end is called after a Span is ended (i.e., the end timestamp is already set). This method is called synchronously within the Span::end API, therefore it should not block or throw an exception.

§Accessing Context

Important: Do not rely on Context::current() in on_end. When on_end is called during span cleanup, Context::current() returns whatever context happens to be active at that moment, which is typically unrelated to the span being ended. Contexts can be activated in any order and are not necessarily hierarchical.

Best Practice: Extract any needed context information in on_start and store it as span attributes. This ensures the information is available in the SpanData passed to on_end.

§Example
ⓘ
impl SpanProcessor for MyProcessor {
    fn on_start(&self, span: &mut Span, cx: &Context) {
        // Extract baggage and store as span attribute
        if let Some(value) = cx.baggage().get("my-key") {
            span.set_attribute(KeyValue::new("my-key", value.to_string()));
        }
    }

    fn on_end(&self, span: SpanData) {
        // Access the attribute stored in on_start
        let my_value = span.attributes.iter()
            .find(|kv| kv.key.as_str() == "my-key");
    }
}
§Filtering completed spans

Warning: Filtering individual spans can produce incomplete or broken traces, such as an exported child whose parent was discarded. This does not coordinate filtering across spans or services. For coordinated decisions based on completed spans, prefer tail-based sampling in the OpenTelemetry Collector or another telemetry pipeline. All spans in a trace must reach the same tail-sampling instance; it cannot recover spans already discarded by SDK sampling or filtering.

If SDK processor filtering fits your requirements and you accept this tradeoff, wrap another processor and delegate only the spans that satisfy your condition. This example uses an attribute, but the condition can use any information in SpanData or other processor state. Register only the wrapper with SdkTracerProvider, since separately registered processors receive spans independently.

use opentelemetry::{Context, Value};
use opentelemetry_sdk::{
    error::OTelSdkResult,
    trace::{Span, SpanData, SpanProcessor},
    Resource,
};
use std::time::Duration;

#[derive(Debug)]
struct FilteringSpanProcessor<P> {
    next: P,
}

impl<P> FilteringSpanProcessor<P> {
    fn new(next: P) -> Self {
        Self { next }
    }
}

impl<P: SpanProcessor> SpanProcessor for FilteringSpanProcessor<P> {
    fn on_start(&self, span: &mut Span, cx: &Context) {
        self.next.on_start(span, cx);
    }

    fn on_end(&self, span: SpanData) {
        let should_drop = span.attributes.iter().any(|attribute| {
            attribute.key.as_str() == "example.drop"
                && attribute.value == Value::Bool(true)
        });

        if !should_drop {
            self.next.on_end(span);
        }
    }

    fn force_flush(&self) -> OTelSdkResult {
        self.next.force_flush()
    }

    fn shutdown_with_timeout(&self, timeout: Duration) -> OTelSdkResult {
        self.next.shutdown_with_timeout(timeout)
    }

    fn set_resource(&mut self, resource: &Resource) {
        self.next.set_resource(resource);
    }
}

TODO - This method should take reference to SpanData

fn force_flush(&self) -> Result<(), OTelSdkError>

Force the spans lying in the cache to be exported.

fn shutdown_with_timeout(&self, timeout: Duration) -> Result<(), OTelSdkError>

Shuts down the processor. Called when SDK is shut down. This is an opportunity for processors to do any cleanup required.

Implementation should make sure shutdown can be called multiple times.

Provided Methods§

fn shutdown(&self) -> Result<(), OTelSdkError>

shutdown the processor with a default timeout.

fn set_resource(&mut self, _resource: &Resource)

Set the resource for the span processor.

Dyn Compatibility§

This trait is dyn compatible.

In older versions of Rust, dyn compatibility was called "object safety".

Implementors§