Rust中用sqlx实现Postgres多线程操作的惯用方法及借检查器问题
问题:Rust中使用结构体方法在Tokio线程执行Postgres插入时遭遇'static生命周期错误
我尝试通过Tokio线程将三种不同数据类型插入Postgres数据库,但遇到了Rust借检查器报错。我曾简化DbMethods结构体,让每个方法接收PgPool引用并克隆连接池,也尝试克隆DbMethods实例,均未解决,出现「dbm需被借用为'static」的错误。目前将数据库操作改为独立函数后线程可正常运行,但希望了解如何以结构体方法的形式实现该功能。
报错代码片段
let connection_pool = PgPool::connect(&connection_string) .await .expect("Failed to connect to Postgres"); let dbm = DbMethods {}; // Make API calls etc.. if let Some(messages) = last_hour.ais_response.ais_latest_responses { // TODO: Handle errors. let split_messages = process_ais_items(messages).unwrap(); // TODO: Create thread for each message type and do DB inserts. let aton_handle = task::spawn(dbm.insert_aton_data( connection_pool.clone(), split_messages.aton_data, &log_id, )); // ... other handles. let _ = tokio::try_join!(aton_handle, static_handle, position_handle); }
方法定义
pub async fn insert_aton_data( &self, db_pool: PgPool, aton_data: Vec<AISAtonData>, log_id: &Uuid, ) -> Result<(), Box<dyn std::error::Error + Send + Sync>> { // let pool = PgPool::connect(&self.connection_string).await?; let tx = db_pool.begin().await?; for data in aton_data { sqlx::query!( "INSERT INTO ais.ais_aton_data ( type_field, message_type, mmsi, msgtime, dimension_a, dimension_b, dimension_c, dimension_d, type_of_aids_to_navigation, latitude, longitude, name, type_of_electronic_fixing_device, log_id ) VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14)", data.type_field, data.message_type, data.mmsi, convert_to_datetime_option(data.msgtime), data.dimension_a, data.dimension_b, data.dimension_c, data.dimension_d, data.type_of_aids_to_navigation, data.latitude, data.longitude, data.name, data.type_of_electronic_fixing_device, log_id ).execute(&db_pool).await?; } tx.commit().await?; Ok(()) }
当前临时解决方案
async fn insert_ais_items(connection_pool: PgPool, log_id: Uuid, last_hour: LastHourAISMessage) -> Result<(), Box<dyn Error>> { if let Some(messages) = last_hour.ais_response.ais_latest_responses { // TODO: Handle errors. let split_messages = process_ais_items(messages).unwrap(); let aton_handle = task::spawn(insert_aton_data( connection_pool.clone(), split_messages.aton_data, log_id.clone(), )); let static_handle = task::spawn(insert_static_data( connection_pool.clone(), split_messages.static_data, log_id.clone(), )); let position_handle = task::spawn(insert_position_data( connection_pool.clone(), split_messages.position_data, log_id.clone(), )); let res = tokio::try_join!(aton_handle, static_handle, position_handle); match res { Ok(..) => { debug!("Threads completed"); } Err(error) => warn!("There was an error in one of the threads: {:?}", error) } } Ok(()) }
问题原因与解决方法
错误原因
Tokio的task::spawn要求传入的Future必须满足'static生命周期,因为任务的执行周期可能超过当前函数作用域。原代码中:
dbm.insert_aton_data方法持有&self引用,而dbm是当前作用域的局部变量,其生命周期无法达到'staticlog_id参数使用了引用&Uuid,同样无法满足'static要求
解决方法
方案1:用Arc包裹DbMethods实例
通过Arc(原子引用计数)共享DbMethods,让每个任务持有独立的Arc引用,满足'static生命周期:
- 修改结构体实例的创建与传递:
use std::sync::Arc; let connection_pool = PgPool::connect(&connection_string) .await .expect("Failed to connect to Postgres"); // 用Arc包裹DbMethods实例 let dbm = Arc::new(DbMethods {}); // ... 其他逻辑 if let Some(messages) = last_hour.ais_response.ais_latest_responses { let split_messages = process_ais_items(messages).unwrap(); // 克隆Arc给每个任务 let aton_dbm = Arc::clone(&dbm); let aton_pool = connection_pool.clone(); let aton_log_id = log_id.clone(); let aton_handle = task::spawn(async move { aton_dbm.insert_aton_data(aton_pool, split_messages.aton_data, aton_log_id).await }); // 其他任务同理,克隆dbm、connection_pool和log_id let static_dbm = Arc::clone(&dbm); let static_pool = connection_pool.clone(); let static_log_id = log_id.clone(); let static_handle = task::spawn(async move { static_dbm.insert_static_data(static_pool, split_messages.static_data, static_log_id).await }); let position_dbm = Arc::clone(&dbm); let position_pool = connection_pool.clone(); let position_log_id = log_id.clone(); let position_handle = task::spawn(async move { position_dbm.insert_position_data(position_pool, split_messages.position_data, position_log_id).await }); let _ = tokio::try_join!(aton_handle, static_handle, position_handle); }
- 修改方法的
log_id参数为拥有所有权:
pub async fn insert_aton_data( &self, db_pool: PgPool, aton_data: Vec<AISAtonData>, log_id: Uuid, // 从&Uuid改为Uuid ) -> Result<(), Box<dyn std::error::Error + Send + Sync>> { // 内部逻辑不变,sqlx可直接使用Uuid类型 let tx = db_pool.begin().await?; for data in aton_data { sqlx::query!( "INSERT INTO ais.ais_aton_data ( type_field, message_type, mmsi, msgtime, dimension_a, dimension_b, dimension_c, dimension_d, type_of_aids_to_navigation, latitude, longitude, name, type_of_electronic_fixing_device, log_id ) VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14)", data.type_field, data.message_type, data.mmsi, convert_to_datetime_option(data.msgtime), data.dimension_a, data.dimension_b, data.dimension_c, data.dimension_d, data.type_of_aids_to_navigation, data.latitude, data.longitude, data.name, data.type_of_electronic_fixing_device, log_id ).execute(&db_pool).await?; } tx.commit().await?; Ok(()) }
方案2:让DbMethods实现Clone
如果DbMethods是无状态的(当前是空结构体),可以直接为其实现Clone,让每个任务持有独立的实例:
- 为结构体添加Clone派生:
#[derive(Clone)] pub struct DbMethods {}
- 修改调用逻辑:
let connection_pool = PgPool::connect(&connection_string) .await .expect("Failed to connect to Postgres"); let dbm = DbMethods {}; // ... 其他逻辑 if let Some(messages) = last_hour.ais_response.ais_latest_responses { let split_messages = process_ais_items(messages).unwrap(); // 克隆DbMethods实例给每个任务 let aton_dbm = dbm.clone(); let aton_pool = connection_pool.clone(); let aton_log_id = log_id.clone(); let aton_handle = task::spawn(async move { aton_dbm.insert_aton_data(aton_pool, split_messages.aton_data, aton_log_id).await }); // 其他任务同理 let static_dbm = dbm.clone(); let static_pool = connection_pool.clone(); let static_log_id = log_id.clone(); let static_handle = task::spawn(async move { static_dbm.insert_static_data(static_pool, split_messages.static_data, static_log_id).await }); let position_dbm = dbm.clone(); let position_pool = connection_pool.clone(); let position_log_id = log_id.clone(); let position_handle = task::spawn(async move { position_dbm.insert_position_data(position_pool, split_messages.position_data, position_log_id).await }); let _ = tokio::try_join!(aton_handle, static_handle, position_handle); }
- 同样需要将方法的
log_id参数改为Uuid(拥有所有权),和方案1一致。
为什么独立函数可以正常运行
独立函数不需要引用self,且参数均通过克隆传递(拥有所有权),自然满足'static生命周期要求,因此可以被task::spawn正常调用。
内容的提问来源于stack exchange,提问作者amsten
相关产品推荐
相关产品推荐

