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

AWS Kinesis Rust SDK与Tokio Runtime兼容问题:异步内同步调用挂起

AWS Kinesis Rust SDK异步与同步桥接进程挂起问题排查

问题现象

尝试将AWS Kinesis Rust SDK的异步功能与同步代码桥接,实现带Stream trait的tokio_stream结构体时,遇到以下异常:

  • 在#[tokio::test]标记的异步测试中,调用第二个包含futures::executor::block_on的同步函数时,进程直接挂起,无任何报错输出
  • 纯异步测试、手动创建Tokio Runtime的同步测试均能正常运行,无异常行为

附带测试代码:

#[cfg(test)]
mod tests {
    use aws_config::meta::region::RegionProviderChain;
    use aws_sdk_kinesis::{Client, Region};
    use aws_config::SdkConfig;
    use aws_sdk_kinesis::error::ListShardsError;
    use aws_sdk_kinesis::output::ListShardsOutput;
    use aws_sdk_kinesis::types::SdkError;
    use futures::executor::block_on;
    use tokio::runtime::Runtime;
    /// 正常运行的纯异步测试
    #[tokio::test]
    async fn testing_kinesis_stream_connection() -> () {
        let region_provider = RegionProviderChain::first_try(Region::new("eu-central-1".to_string()));
        let sdk_config =
            aws_config::from_env()
                .region(region_provider)
                .load().await;
        let client = Client::new(&sdk_config);
        let shard_request = client.list_shards().stream_name("datafusion-stream".to_string());
        let shards = shard_request.send().await.unwrap();
        println!("{:?}", shards);
    }
    /// 正常运行的手动Runtime同步测试
    #[test]
    fn testing_kinesis_stream_connection_without_tokio() -> () {
        let runtime = tokio::runtime::Builder::new_current_thread().enable_all().build().unwrap();
        let region_provider = RegionProviderChain::first_try(Region::new("eu-central-1".to_string()));
        let sdk_config =
            runtime.block_on(aws_config::from_env()
                .region(region_provider)
                .load());
        let client = Client::new(&sdk_config);
        let shard_request = client.list_shards().stream_name("datafusion-stream".to_string());
        let shards = runtime.block_on(shard_request.send()).unwrap();
        println!("{:?}", shards);
    }

    fn this_needs_to_be_sync_1() -> SdkConfig {
        let region_provider = RegionProviderChain::first_try(Region::new("eu-central-1".to_string()));
        let sdk_config =
            futures::executor::block_on(aws_config::from_env()
                .region(region_provider)
                .load());
        sdk_config
    }

    fn this_needs_to_be_sync_2(client: Client) -> ListShardsOutput {
        let shard_request = client.list_shards().stream_name("datafusion-stream".to_string());
        let shards = futures::executor::block_on(shard_request.send()).unwrap();
        shards
    }
    /// 挂起的测试:第二个block_on处停滞
    #[tokio::test]
    async fn testing_kinesis_stream_connection_sync_inside_async() -> () {
        // 第一个block_on正常执行
        let sdk_config = this_needs_to_be_sync_1();
        let client = Client::new(&sdk_config);
        // 此处挂起
        let shards = this_needs_to_be_sync_2(client);
        println!("{:?}", shards);
    }

}

问题根源

核心是嵌套阻塞与Tokio Runtime上下文冲突:

  • #[tokio::test]默认创建的Tokio Runtime,其异步任务运行在专属的IO线程池中
  • 当在Tokio异步任务中调用futures::executor::block_on时,第二个block_on会阻塞当前Tokio线程,等待内部future完成
  • 但AWS Kinesis SDK的send() future依赖Tokio Runtime的IO驱动来执行网络请求,此时被阻塞的Tokio线程无法处理IO操作,导致future永远无法完成,形成死锁
  • 手动创建Runtime的测试能正常运行,是因为所有异步操作都在同一个Runtime的block_on下,上下文统一;纯异步测试全程使用Tokio的异步调度,无阻塞调用,因此没有冲突

解决方法

方法1:统一使用Tokio Runtime处理所有阻塞

将同步函数中的futures::executor::block_on替换为Tokio Runtime的block_on,确保所有异步操作在同一个上下文执行:

use once_cell::sync::OnceCell;
use tokio::runtime::Runtime;

// 全局初始化Tokio Runtime,确保所有同步阻塞都用同一个Runtime
fn get_global_tokio_runtime() -> &'static Runtime {
    static RUNTIME: OnceCell<Runtime> = OnceCell::new();
    RUNTIME.get_or_init(|| {
        tokio::runtime::Builder::new_multi_thread()
            .enable_all()
            .build()
            .unwrap()
    })
}

fn this_needs_to_be_sync_1() -> SdkConfig {
    let region_provider = RegionProviderChain::first_try(Region::new("eu-central-1".to_string()));
    // 使用全局Tokio Runtime的block_on
    let sdk_config = get_global_tokio_runtime().block_on(aws_config::from_env()
        .region(region_provider)
        .load());
    sdk_config
}

fn this_needs_to_be_sync_2(client: Client) -> ListShardsOutput {
    let shard_request = client.list_shards().stream_name("datafusion-stream".to_string());
    // 同样使用全局Tokio Runtime的block_on
    let shards = get_global_tokio_runtime().block_on(shard_request.send()).unwrap();
    shards
}

方法2:消除嵌套阻塞,改用全异步调用

如果同步函数的限制可以放宽,直接将其改为异步函数,在测试中用await调用,避免任何阻塞:

// 改为异步函数
async fn this_should_be_async_1() -> SdkConfig {
    let region_provider = RegionProviderChain::first_try(Region::new("eu-central-1".to_string()));
    aws_config::from_env()
        .region(region_provider)
        .load().await
}

async fn this_should_be_async_2(client: Client) -> ListShardsOutput {
    let shard_request = client.list_shards().stream_name("datafusion-stream".to_string());
    shard_request.send().await.unwrap()
}

#[tokio::test]
async fn testing_kinesis_stream_connection_sync_inside_async_fixed() -> () {
    let sdk_config = this_should_be_async_1().await;
    let client = Client::new(&sdk_config);
    let shards = this_should_be_async_2(client).await;
    println!("{:?}", shards);
}

方法3:用spawn_blocking隔离同步阻塞

如果必须保留同步函数,在异步测试中用Tokio的spawn_blocking将同步调用放到独立线程执行,避免阻塞Tokio的IO线程:

#[tokio::test]
async fn testing_kinesis_stream_connection_sync_inside_async_fixed() -> () {
    // 把同步函数放到spawn_blocking中执行
    let sdk_config = tokio::task::spawn_blocking(|| this_needs_to_be_sync_1()).await.unwrap();
    let client = Client::new(&sdk_config);
    let shards = tokio::task::spawn_blocking(|| this_needs_to_be_sync_2(client)).await.unwrap();
    println!("{:?}", shards);
}

内容的提问来源于stack exchange,提问作者Metehan Yıldırım

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.13 14:40:28