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
相关产品推荐
相关产品推荐

