Axum SSE端点每秒调用异步函数报错问题求助
解决Axum SSE端点异步定时查询数据库的问题
你的核心问题是stream::repeat_with不支持异步闭包——它生成的是同步迭代器流,无法直接处理异步操作。改成async闭包后,流的元素变成了Future,需要额外处理才能展开成实际值。下面是两种可行的修复方案:
方案一:用IntervalStream实现定时异步任务
这是最直观的方式,用Tokio的Interval控制每秒触发一次,绑定异步查询逻辑:
use tokio_stream::wrappers::IntervalStream; use tokio::time::{interval, Duration}; pub async fn payment_sse_handler( State(state): State<AppState>, Extension(current_user): Extension<CurrentUser>, ) -> Sse<impl Stream<Item = Result<Event, Infallible>>> { // 创建每秒触发一次的定时器 let interval = interval(Duration::from_secs(1)); let stream = IntervalStream::new(interval) // 对每个定时触发事件,执行异步数据库查询 .then(move |_| { // 克隆资源供闭包复用,Pool本身是Arc包装,克隆开销极小 let db_pool = state.db_pool.clone(); let user_id = current_user.user_id.clone(); async move { let event = get_balance(&db_pool, user_id).await; Ok(event) } }); Sse::new(stream).keep_alive(KeepAlive::default()) } // 保留原有的get_balance函数不变 pub async fn get_balance (db_pool: &Pool<Postgres>, user_id: String) -> Event { let transaction_history = sqlx::query_as!( Transaction, r#"select transaction_id, status, currency, amount, user_id, created_at, processing_fee from transaction_history WHERE user_id=$1 ORDER BY created_at LIMIT 5"#, user_id ) .fetch_all(db_pool) .await; match transaction_history { Ok(transaction_history) => { serde_json::to_string(&transaction_history) .map(|s| Event::default().data(s)) .unwrap_or_else(|_| Event::default().data("Error")) } Err(_) => Event::default().data("Error"), } }
方案二:用repeat_with配合then展开Future
如果坚持用repeat_with,可以先生成Future流,再通过then操作展开获取实际值:
use futures::stream::{self, StreamExt}; pub async fn payment_sse_handler( State(state): State<AppState>, Extension(current_user): Extension<CurrentUser>, ) -> Sse<impl Stream<Item = Result<Event, Infallible>>> { let stream = stream::repeat_with(move || { let db_pool = state.db_pool.clone(); let user_id = current_user.user_id.clone(); // 返回Future而非直接await async move { get_balance(&db_pool, user_id).await } }) // 控制每秒执行一次查询 .throttle(Duration::from_secs(1)) // 展开Future,获取实际的Event结果 .then(|fut| fut) .map(Ok); Sse::new(stream).keep_alive(KeepAlive::default()) }
关键说明
repeat_with生成的是Stream<Item=F>(F为Future类型),必须通过then或buffer_unordered等操作,将Future转换成实际输出值才能继续处理。IntervalStream更适配定时任务场景,基于Tokio原生定时器,时间控制更准确,代码可读性更高。- 克隆
db_pool和user_id是因为闭包会被多次调用,需确保每次调用都有独立可访问的资源。
内容的提问来源于stack exchange,提问作者Joanthan Ahrenkiel-Frellsen
相关产品推荐
相关产品推荐

