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

Rust项目中如何为异步函数正确添加tracing span?

问题

在Rust项目中为latest_block_header函数添加tracing追踪时出现异常:程序运行后,追踪数据发送到Jaeger后,无法看到该函数对应的testest和latest_block_header in eth adapter span,仅能看到PollingBlockIngestor::do_poll、latest_block等其他span。

相关代码如下:

chain/ethereum/src/ingestor.rs

#[tracing::instrument(skip(self), name = "PollingBlockIngestor::do_poll")]
async fn do_poll(&self) -> Result<(), IngestorError> {

    // Get chain head ptr from store
    let head_block_ptr_opt = self.chain_store.cheap_clone().chain_head_ptr().await?;

    // To check if there is a new block or not, fetch only the block header since that's cheaper
    // than the full block. This is worthwhile because most of the time there won't be a new
    // block, as we expect the poll interval to be much shorter than the block time.
    info_span!("latest_block");
    let latest_block = self.latest_block().await?;

    if let Some(head_block) = head_block_ptr_opt.as_ref() {
        // If latest block matches head block in store, nothing needs to be done
        if &latest_block == head_block {
            return Ok(());
        }

        if latest_block.number < head_block.number {
            // An ingestor might wait or move forward, but it never
            // wavers and goes back. More seriously, this keeps us from
            // later trying to ingest a block with the same number again
            warn!(self.logger,
                "Provider went backwards - ignoring this latest block";
                "current_block_head" => head_block.number,
                "latest_block_head" => latest_block.number);
            return Ok(());
        }
    }

    // Compare latest block with head ptr, alert user if far behind
    match head_block_ptr_opt {
        None => {
            info!(
                self.logger,
                "Downloading latest blocks from Ethereum, this may take a few minutes..."
            );
        }
        Some(head_block_ptr) => {
            let latest_number = latest_block.number;
            let head_number = head_block_ptr.number;
            let distance = latest_number - head_number;
            let blocks_needed = (distance).min(self.ancestor_count);
            let code = if distance >= 15 {
                LogCode::BlockIngestionLagging
            } else {
                LogCode::BlockIngestionStatus
            };
            if distance > 0 {
                info!(
                    self.logger,
                    "Syncing {} blocks from Ethereum",
                    blocks_needed;
                    "current_block_head" => head_number,
                    "latest_block_head" => latest_number,
                    "blocks_behind" => distance,
                    "blocks_needed" => blocks_needed,
                    "code" => code,
                );
            }
        }
    }

    let mut missing_block_hash = self.ingest_block(&latest_block.hash).await?;

    
    while let Some(hash) = missing_block_hash {
        missing_block_hash = self.ingest_block(&hash).await?;
    }
    Ok(())
}

async fn latest_block(&self) -> Result<BlockPtr, IngestorError> {
    info_span!("latest_block_header");
    self.eth_adapter
        .latest_block_header(&self.logger)
        .compat()
        .await
        .map(|block| block.into())
}

chain/ethereum/src/ethereum_adapter.rs

impl EthereumAdapterTrait for EthereumAdapter {

    #[tracing::instrument(skip_all, name = "testest")]
    fn latest_block_header(
        &self,
        logger: &Logger,
    ) -> Box<dyn Future<Item = web3::types::Block<H256>, Error = IngestorError> + Send> {
        let s = info_span!("latest_block_header in eth adapter");
        let web3 = self.web3.clone();
        Box::new(
            retry("eth_getBlockByNumber(latest) no txs RPC call", logger)
                .no_limit()
                .timeout_secs(ENV_VARS.json_rpc_timeout.as_secs())
                .run(move || {
                    let web3 = web3.cheap_clone();
                    async move {
                        let block_opt = web3
                            .eth()
                            .block(Web3BlockNumber::Latest.into())
                            .await
                            .map_err(|e| {
                                anyhow!("could not get latest block from Ethereum: {}", e)
                            })?;

                        block_opt
                            .ok_or_else(|| anyhow!("no latest block returned from Ethereum").into())
                    }
                })
                .map_err(move |e| {
                    e.into_inner().unwrap_or_else(move || {
                        anyhow!("Ethereum node took too long to return latest block").into()
                    })
                })
                .boxed()
                .compat(),
        )
    }
// lots of other functions
}

曾尝试修改async代码块,对future调用instrument方法,但依然无效:

let web3 = web3.cheap_clone();
let s = info_span!("latest_block_header in eth adapter");
async move {
    let block_opt = web3
        .eth()
        .block(Web3BlockNumber::Latest.into())
        .instrument(s)
        .await
        .map_err(|e| {
            anyhow!("could not get latest block from Ethereum: {}", e)
        })?;

    block_opt
        .ok_or_else(|| anyhow!("no latest block returned from Ethereum").into())
}

解决方案

1. 修复#[tracing::instrument]的无效问题

#[tracing::instrument]宏对返回Box<dyn Future>的同步函数无法自动生效——它只会在函数调用时创建span,但不会将span附加到返回的Future上。如果无法修改trait定义将函数改为async fn,需手动创建span并绑定到返回的Future:

fn latest_block_header(
    &self,
    logger: &Logger,
) -> Box<dyn Future<Item = web3::types::Block<H256>, Error = IngestorError> + Send> {
    // 创建目标span
    let span = tracing::info_span!("testest");
    let web3 = self.web3.clone();
    
    // 构建原始future
    let future = retry("eth_getBlockByNumber(latest) no txs RPC call", logger)
        .no_limit()
        .timeout_secs(ENV_VARS.json_rpc_timeout.as_secs())
        .run(move || {
            let web3 = web3.cheap_clone();
            async move {
                let block_opt = web3
                    .eth()
                    .block(Web3BlockNumber::Latest.into())
                    .await
                    .map_err(|e| {
                        anyhow!("could not get latest block from Ethereum: {}", e)
                    })?;

                block_opt
                    .ok_or_else(|| anyhow!("no latest block returned from Ethereum").into())
            }
        })
        .map_err(move |e| {
            e.into_inner().unwrap_or_else(move || {
                anyhow!("Ethereum node took too long to return latest block").into()
            })
        })
        .boxed()
        .compat();
    
    // 将span绑定到整个future,确保执行期间span处于激活状态
    Box::new(future.instrument(span))
}

2. 修复手动创建的latest_block_header in eth adapter span

之前的代码仅创建了span但未绑定到对应的异步任务,需将其附加到retry内部的async块Future:

.run(move || {
    let web3 = web3.cheap_clone();
    let inner_span = tracing::info_span!("latest_block_header in eth adapter");
    // 将async块的Future用span包裹
    async move {
        let block_opt = web3
            .eth()
            .block(Web3BlockNumber::Latest.into())
            .await
            .map_err(|e| {
                anyhow!("could not get latest block from Ethereum: {}", e)
            })?;

        block_opt
            .ok_or_else(|| anyhow!("no latest block returned from Ethereum").into())
    }
    .instrument(inner_span)
})

3. 修复latest_block函数中的无效span

ingestor.rs中的latest_block函数仅创建了span但未激活,需将其绑定到后续的Future:

async fn latest_block(&self) -> Result<BlockPtr, IngestorError> {
    let span = info_span!("latest_block_header");
    // 将后续操作的Future用span包裹
    self.eth_adapter
        .latest_block_header(&self.logger)
        .compat()
        .instrument(span)
        .await
        .map(|block| block.into())
}

核心原因总结

  • #[tracing::instrument]不会自动处理返回Box<dyn Future>的同步函数,必须手动将span附加到Future。
  • 手动创建的span如果没有通过instrument()绑定到Future,或者未调用span.enter()激活,不会被tracing系统记录。
  • 所有异步操作的span必须和对应的Future绑定,才能在Future执行期间保留追踪上下文。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 12:35:54