You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在Rust自定义Tokio-Tracing订阅器中避免Hyper无限日志循环

问题:仅抑制订阅器内部依赖库的追踪事件,保留应用层同类日志

说明

我已查阅《在tokio tracing订阅器中抑制外部库日志》和《如何关闭其他 crate 发出的追踪事件?》,但它们无法解决我的问题。

背景

我正在开发一个Tokio tracing订阅器,用于捕获由tracing instrumentation的库生成的日志,并通过hyper crate导出到后端服务。但hyper同样基于Tokio tracing进行了 instrumentation,导致其生成的日志被我的订阅器捕获并导出,进而引发无限日志循环。

请注意,集成该订阅器的应用也可能使用hyper crate。我仅希望抑制订阅器实现内部使用hyper所产生的追踪事件,而非应用层使用hyper产生的事件。如链接问题中建议的添加hyper过滤器会抑制所有hyper事件,无论其使用场景。以下是简化的可运行代码及问题细节:

/*
hyper = { version = "0.14.7", features = [ "full" ] }
tokio = { version = "1.33.0", features = ["full"] }
tracing = "0.1.25"
tracing-core = "0.1"
pin-project = "1.1.3"
tracing-subscriber = "0.3.17"
*/

use pin_project::pin_project;
use std::future::Future;
use std::pin::Pin;
use std::task::{Context, Poll};
use tracing::info;
use tracing::Event;
use tracing::Metadata;
use tracing::Subscriber;
use tracing_core::span::Id;
use tracing_core::span::Record;

// Define a task-local variable for suppression
tokio::task_local! {
    static SUPPRESSED: bool;
}

struct SimpleSubscriber;

impl Subscriber for SimpleSubscriber {
    fn enabled(&self, _metadata: &Metadata<'_>) -> bool {
        !is_logging_suppressed()
    }

    fn new_span(&self, _: &tracing::span::Attributes<'_>) -> Id {
        Id::from_u64(10)
    }

    fn event(&self, event: &Event<'_>) {
        if is_logging_suppressed() {
            return;
        }

        // Extract the required metadata before moving into the async context
        let target = event.metadata().target().to_string();
        let name = event.metadata().name().to_string();
        let level = event.metadata().level().clone();

        let suppressed_future = SUPPRESSED.scope(true, async move {
            // the actual logging event, the hyper tracing event generated here should be suppressed
            hyper_wrapper::make_hyper_call("mock_upload_event", &target, &name, &level).await;
        });

        let suppress_logging_future = SuppressLogging {
            inner: suppressed_future,
        };

        // Use tokio::spawn instead of spawn_local
        tokio::task::spawn(suppress_logging_future);
    }

    fn record(&self, _: &Id, _: &Record<'_>) {}

    fn record_follows_from(&self, _: &Id, _: &Id) {}

    fn enter(&self, _: &Id) {}

    fn exit(&self, _: &Id) {}
}

#[pin_project]
struct SuppressLogging<F> {
    #[pin]
    inner: F,
}

impl<F: Future> Future for SuppressLogging<F> {
    type Output = F::Output;

    fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
        let this = self.project();
        this.inner.poll(cx)
    }
}

fn is_logging_suppressed() -> bool {
    SUPPRESSED.try_with(|val| *val).unwrap_or(false)
}

#[tokio::main]
async fn main() {
    let subscriber = SimpleSubscriber;
    tracing::subscriber::set_global_default(subscriber)
        .expect("Setting default subscriber failed.");

    // Make a small Hyper call in the main method
    lib::some_func().await;
}

// Define the module within the same file
mod lib {
    use tracing::info;
    use tracing_core::Level;

    pub async fn some_func() {
        // this should be logged
        info!(name: "this-is-also-logged", "test");

        // the tracing event generated by hyper crate should not be suppressed.
        super::hyper_wrapper::make_hyper_call("main method", "main", "startup", &Level::INFO).await;
    }
}

mod hyper_wrapper {
    use hyper::Client;
    use hyper::Uri;
    use tracing_core::Level;

    pub async fn make_hyper_call(context: &str, target: &str, name: &str, level: &Level) {
        let client = Client::new();
        let uri = "http://www.google.com".parse::<Uri>().unwrap();
        let response = client.get(uri).await;
        match response {
            Ok(_response) => {
                //let _body = hyper::body::to_bytes(response).await.unwrap();
                println!(
                    "{} - [{}] - {} - {} ",
                    context,
                    level,
                    target,
                    name,
                );
            }
            Err(e) => {
                println!("{} - [{}] - {} - {} - Error: {:?}", context, level, target, name, e);
            }
        }
    }
}

问题说明

  • 我的Tokio tracing订阅器通过hyper将日志导出到后端,但hyper自身的追踪事件会被订阅器捕获,引发无限循环;
  • 直接过滤所有hyper日志会影响应用层或其他依赖库的hyper日志;
  • 我尝试使用tokio::task_local设置抑制标记,但该标记无法在hyper内部的异步调用链中传播,导致抑制失效。

核心问题

如何仅过滤订阅器实现内部hyper(及reqwest、tonic等同类crate)产生的日志,同时不影响应用层或其他地方的hyper日志?

示例场景说明

  • SimpleSubscriber是我将发布为crate的订阅器;
  • main函数属于集成该crate的应用,同时使用lib模块;
  • lib模块的日志及其中使用hyper产生的日志应被订阅器导出;
  • 订阅器内部通过event()调用make_hyper_call时产生的hyper日志需被抑制,避免循环;
  • 尝试的tokio::task_local方案因标记无法传播而失效。

解决方案

要解决这个问题,核心是区分订阅器内部发起的异步任务和应用层的任务,确保抑制逻辑能准确覆盖内部操作,同时不干扰外部任务。以下是几种可行方案:

1. 基于tracing span的上下文判断

在订阅器内部发起导出操作时,创建一个专属标记span,然后在enabled方法中检查当前是否处于该span内,同时判断事件是否来自需要抑制的crate:

impl Subscriber for SimpleSubscriber {
    fn enabled(&self, metadata: &Metadata<'_>) -> bool {
        // 检查当前是否处于订阅器导出的span中
        let is_in_export_span = tracing::Span::current().name() == "subscriber_export";
        // 定义需要抑制的目标crate
        let is_suppressed_crate = metadata.target().starts_with("hyper") 
            || metadata.target().starts_with("reqwest")
            || metadata.target().starts_with("tonic");

        // 仅当处于导出span且是目标crate时才抑制
        !(is_in_export_span && is_suppressed_crate)
    }

    fn event(&self, event: &Event<'_>) {
        if !self.enabled(event.metadata()) {
            return;
        }

        let target = event.metadata().target().to_string();
        let name = event.metadata().name().to_string();
        let level = event.metadata().level().clone();

        tokio::spawn(async move {
            // 创建专属span标记导出操作,span会自动在异步链中传播
            let export_span = tracing::info_span!("subscriber_export");
            let _guard = export_span.enter();
            
            hyper_wrapper::make_hyper_call("mock_upload_event", &target, &name, &level).await;
        });
    }

    // 其他方法保持不变...
}

这种方式无需依赖任务本地变量,利用tracing自身的span上下文传播特性,可靠性更高。

2. 修复任务本地变量的传播问题

如果坚持使用tokio::task_local,需要确保所有内部异步操作都在标记的scope内执行,同时去掉多余的SuppressLogging包装器(它不会帮助传播上下文):

fn event(&self, event: &Event<'_>) {
    if is_logging_suppressed() {
        return;
    }

    let target = event.metadata().target().to_string();
    let name = event.metadata().name().to_string();
    let level = event.metadata().level().clone();

    // 直接在spawn的任务中设置scope,确保所有内部异步调用都继承该标记
    tokio::spawn(async move {
        SUPPRESSED.scope(true, async move {
            hyper_wrapper::make_hyper_call("mock_upload_event", &target, &name, &level).await;
        }).await;
    });
}

如果需要使用tokio::spawn_local(保留线程本地上下文),需确保任务是!Send类型,或者通过tokio::sync::oneshot传递结果。

3. 结合任务ID标记

可以在订阅器内部生成任务时记录任务ID,然后在enabled方法中检查当前任务ID是否属于订阅器内部任务:

use std::sync::atomic::{AtomicU64, Ordering};

static INTERNAL_TASK_ID: AtomicU64 = AtomicU64::new(0);

impl Subscriber for SimpleSubscriber {
    fn enabled(&self, metadata: &Metadata<'_>) -> bool {
        let is_suppressed_crate = metadata.target().starts_with("hyper") 
            || metadata.target().starts_with("reqwest")
            || metadata.target().starts_with("tonic");
        
        if !is_suppressed_crate {
            return true;
        }

        // 检查当前任务是否是订阅器内部任务
        let current_task_id = tokio::task::Id::current().as_u64();
        let internal_task_id = INTERNAL_TASK_ID.load(Ordering::Relaxed);
        current_task_id != internal_task_id
    }

    fn event(&self, event: &Event<'_>) {
        if !self.enabled(event.metadata()) {
            return;
        }

        let target = event.metadata().target().to_string();
        let name = event.metadata().name().to_string();
        let level = event.metadata().level().clone();

        let task = tokio::spawn(async move {
            hyper_wrapper::make_hyper_call("mock_upload_event", &target, &name, &level).await;
        });

        // 记录内部任务ID
        INTERNAL_TASK_ID.store(task.id().as_u64(), Ordering::Relaxed);
    }

    // 其他方法保持不变...
}

这种方式适合简单场景,但需要注意任务ID的复用问题,可结合任务完成后的清理逻辑优化。


内容的提问来源于stack exchange,提问作者LKB

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.23 18:43:09