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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.07 04:45:54