improve code for opentelemetry

This commit is contained in:
houseme
2025-04-26 17:26:27 +08:00
parent ed4a9db9a0
commit 05f1412323
2 changed files with 37 additions and 53 deletions
+30 -53
View File
@@ -3,7 +3,7 @@ use crate::utils::get_local_ip_with_default;
use crate::OtelConfig; use crate::OtelConfig;
use opentelemetry::trace::TracerProvider; use opentelemetry::trace::TracerProvider;
use opentelemetry::{global, KeyValue}; use opentelemetry::{global, KeyValue};
use opentelemetry_appender_tracing::layer; use opentelemetry_appender_tracing::layer::OpenTelemetryTracingBridge;
use opentelemetry_otlp::WithExportConfig; use opentelemetry_otlp::WithExportConfig;
use opentelemetry_sdk::logs::SdkLoggerProvider; use opentelemetry_sdk::logs::SdkLoggerProvider;
use opentelemetry_sdk::{ use opentelemetry_sdk::{
@@ -202,39 +202,6 @@ pub fn init_telemetry(config: &OtelConfig) -> OtelGuard {
// configuring tracing // configuring tracing
{ {
// optimize filter configuration
let otel_layer = {
let filter_otel = match logger_level {
"trace" | "debug" => {
info!("OpenTelemetry tracing initialized with level: {}", logger_level);
EnvFilter::new(logger_level)
}
_ => {
let mut filter = EnvFilter::new(logger_level);
// use smallvec to avoid heap allocation
let directives: SmallVec<[&str; 5]> = smallvec::smallvec!["hyper", "tonic", "h2", "reqwest", "tower"];
for directive in directives {
filter = filter.add_directive(format!("{}=off", directive).parse().unwrap());
}
filter
}
};
layer::OpenTelemetryTracingBridge::new(&logger_provider).with_filter(filter_otel)
};
let tracer = tracer_provider.tracer(Cow::Borrowed(service_name).to_string());
// Configure registry to avoid repeated calls to filter methods
let level_filter = switch_level(logger_level);
let registry = tracing_subscriber::registry()
.with(level_filter)
.with(OpenTelemetryLayer::new(tracer))
.with(MetricsLayer::new(meter_provider.clone()))
.with(otel_layer)
.with(EnvFilter::try_from_default_env().unwrap_or_else(|_| EnvFilter::new(logger_level)));
// configure the formatting layer // configure the formatting layer
let enable_color = std::io::stdout().is_terminal(); let enable_color = std::io::stdout().is_terminal();
let fmt_layer = tracing_subscriber::fmt::layer() let fmt_layer = tracing_subscriber::fmt::layer()
@@ -243,16 +210,25 @@ pub fn init_telemetry(config: &OtelConfig) -> OtelGuard {
.with_thread_names(true) .with_thread_names(true)
.with_thread_ids(true) .with_thread_ids(true)
.with_file(true) .with_file(true)
.with_line_number(true) .with_line_number(true);
.with_filter(
EnvFilter::new(logger_level).add_directive(
format!("opentelemetry={}", if endpoint.is_empty() { logger_level } else { "off" })
.parse()
.unwrap(),
),
);
registry.with(ErrorLayer::default()).with(fmt_layer).init(); let filter = build_env_filter(logger_level, None);
let tracer = tracer_provider.tracer(Cow::Borrowed(service_name).to_string());
let otel_filter = build_env_filter(logger_level, None);
let otel_layer = OpenTelemetryTracingBridge::new(&logger_provider).with_filter(otel_filter);
// Configure registry to avoid repeated calls to filter methods
tracing_subscriber::registry()
.with(filter)
.with(fmt_layer)
.with(ErrorLayer::default())
.with(otel_layer)
.with(MetricsLayer::new(meter_provider.clone()))
.with(OpenTelemetryLayer::new(tracer))
.with(ErrorLayer::default())
.init();
if !endpoint.is_empty() { if !endpoint.is_empty() {
info!( info!(
@@ -269,15 +245,16 @@ pub fn init_telemetry(config: &OtelConfig) -> OtelGuard {
} }
} }
/// Switch log level fn build_env_filter(logger_level: &str, default_level: Option<&str>) -> EnvFilter {
fn switch_level(logger_level: &str) -> tracing_subscriber::filter::LevelFilter { let level = default_level.unwrap_or(logger_level);
use tracing_subscriber::filter::LevelFilter; let mut filter = EnvFilter::try_from_default_env().unwrap_or_else(|_| EnvFilter::new(level));
match logger_level {
"error" => LevelFilter::ERROR, if !matches!(logger_level, "trace" | "debug") {
"warn" => LevelFilter::WARN, let directives: SmallVec<[&str; 5]> = smallvec::smallvec!["hyper", "tonic", "h2", "reqwest", "tower"];
"info" => LevelFilter::INFO, for directive in directives {
"debug" => LevelFilter::DEBUG, filter = filter.add_directive(format!("{}=off", directive).parse().unwrap());
"trace" => LevelFilter::TRACE, }
_ => LevelFilter::OFF,
} }
filter
} }
+7
View File
@@ -335,6 +335,12 @@ async fn run(opt: config::Opt) -> Result<()> {
span span
}) })
.on_request(|request: &HttpRequest<_>, _span: &Span| { .on_request(|request: &HttpRequest<_>, _span: &Span| {
info!(
counter.rustfs_api_requests_total = 1_u64,
key_request_method = %request.method().to_string(),
key_request_uri_path = %request.uri().path().to_owned(),
"handle request api total",
);
debug!("http started method: {}, url path: {}", request.method(), request.uri().path()) debug!("http started method: {}, url path: {}", request.method(), request.uri().path())
}) })
.on_response(|response: &Response<_>, latency: Duration, _span: &Span| { .on_response(|response: &Response<_>, latency: Duration, _span: &Span| {
@@ -342,6 +348,7 @@ async fn run(opt: config::Opt) -> Result<()> {
debug!("http response generated in {:?}", latency) debug!("http response generated in {:?}", latency)
}) })
.on_body_chunk(|chunk: &Bytes, latency: Duration, _span: &Span| { .on_body_chunk(|chunk: &Bytes, latency: Duration, _span: &Span| {
info!(histogram.request.body.len = chunk.len(), "histogram request body lenght",);
debug!("http body sending {} bytes in {:?}", chunk.len(), latency) debug!("http body sending {} bytes in {:?}", chunk.len(), latency)
}) })
.on_eos(|_trailers: Option<&HeaderMap>, stream_duration: Duration, _span: &Span| { .on_eos(|_trailers: Option<&HeaderMap>, stream_duration: Duration, _span: &Span| {