异步Rust中使用Redis连接管理器时遇到future无法在线程间安全共享的错误求助
我来帮你分析一下这个问题,你遇到的核心矛盾是:外部Trait要求read_from_db返回的Future必须同时满足Send + Sync,但Redis库返回的mget操作Future只实现了Send,没有实现Sync,再加上你的连接管理器本身可能也不满足Sync,最终导致生成的异步Future无法被线程安全共享。
先再明确一下你的场景和报错:
场景代码
外部Trait定义:
/// Address space. type Address; /// Values space. type LocalValue; /// Memory error. type Error: Send + Sync + std::error::Error; /// Address, LocalValue and Error are defined /// Reads the words from the given addresses. fn read_from_db( &self, a: Vec<Self::Address>, ) -> impl Send + Sync + Future<Output = Result<Vec<Option<Self::LocalValue>>, Self::Error>>;
你的实现代码:
async fn read_from_db( &self, addresses: Vec<Address>, ) -> Result<Vec<Option<Self::LocalValue>>, Self::Error> { let refs: Vec<&Address> = addresses.iter().collect(); let value = self.clone().connection.mget::<_, Vec<_>>(&refs).await?; Ok(value) }
报错信息
error: future cannot be shared between threads safely | 119 | / async fn batch_read( 120 | | &self, 121 | | addresses: Vec<Address>, 122 | | ) -> Result<Vec<Option<Self::LocalValue>>, Self::Error> { | |_____________________________________________________^ future returned by `read_from_db` is not `Sync` note: future is not `Sync` as it awaits another future which is not `Sync` let value = self.clone().connection.mget::<_, Vec<_>>(&refs).await?; | ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ await occurs here on type `Pin<Box<dyn Future<Output = Result<Vec<Option</* some type */>>, RedisError>> + Send>>`, which is not `Sync` note: required by a bound in /* some trait */ | 70 | ) -> impl Send + Sync + Future<Output = Result<Vec<Option</* some type */>>, Self::Error>>; | ^^^^ required by this bound in /* some trait */
解决方案
方案1:用线程安全的共享容器包装Redis连接管理器
问题的关键在于你的Redis连接管理器(比如redis::aio::ConnectionManager)本身只实现了Send,没有实现Sync,导致持有它的Future无法满足Sync约束。我们可以用Arc<Mutex>把连接管理器包装起来,让它能被安全地跨线程共享:
首先,修改你的结构体定义(以Tokio生态为例):
use std::sync::Arc; use tokio::sync::Mutex; use redis::aio::ConnectionManager; struct YourStruct { // 用Arc+Mutex包装连接管理器,实现线程安全共享 connection: Arc<Mutex<ConnectionManager>>, // 其他字段... }
然后修改read_from_db的实现:
async fn read_from_db( &self, addresses: Vec<Address>, ) -> Result<Vec<Option<Self::LocalValue>>, Self::Error> { let refs: Vec<&Address> = addresses.iter().collect(); // 获取连接管理器的异步锁,这一步的Future是满足Sync的 let mut conn = self.connection.lock().await; // 执行Redis操作 let value = conn.mget::<_, Vec<_>>(&refs).await?; Ok(value) }
这样处理后,整个返回的Future会满足Send + Sync约束:Arc是线程安全的共享指针,tokio::sync::Mutex的锁Future实现了Sync,而且锁机制确保同一时间只有一个线程访问连接,避免了并发冲突。
如果你用的是通用异步 runtime(不是Tokio),可以换成async-lock库的Mutex,用法基本一致,记得先在Cargo.toml添加依赖:
async-lock = "3.0"
方案2:显式包装Future为Sync类型
如果你不想修改连接管理器的包装方式,也可以把整个异步逻辑包装成一个满足Send + Sync的Trait对象:
use futures::future::BoxFuture; fn read_from_db( &self, addresses: Vec<Address>, ) -> BoxFuture<'_, Result<Vec<Option<Self::LocalValue>>, Self::Error>> { let conn = self.connection.clone(); // 用async move捕获conn,然后包装成同时满足Send+Sync的BoxFuture Box::pin(async move { let refs: Vec<&Address> = addresses.iter().collect(); let value = conn.mget::<_, Vec<_>>(&refs).await?; Ok(value) }) as Box<dyn Future<Output = _> + Send + Sync + '_> }
不过这个方案的前提是conn本身是Send的(Redis连接管理器满足这一点),而且这种方式不如第一种方案直观,推荐优先用方案1。
为什么之前的尝试没生效?
你之前试的锁、复制可能没找对位置:比如如果只是复制self,但连接管理器本身还是非Sync的,那复制后的对象依然不满足Sync约束,需要用Arc<Mutex>这种线程安全的容器来包装核心的非Sync资源,才能从根源上解决问题。
备注:内容来源于stack exchange,提问作者Gattouso

