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

如何用tokio::time::timeout包装reqwest请求?类型错误求解

解决Rust中reqwest异步请求Pushshift API的超时与类型匹配问题

问题描述

我用reqwest实现了轮询Pushshift Reddit API的异步逻辑,但长时间运行后API偶尔会出现无限等待的情况。计划通过tokio::time::timeout为请求添加超时限制,却遇到类型匹配错误:最初定义的resp变量类型为String,后续尝试为其赋值send()返回的Future时,编译器提示“expected struct std::string::String, found opaque type”等错误。

我困惑的是,tokio::time::timeout要求传入Future,我传入了send()返回的Future,为何编译器仍期望String类型?虽尝试取消变量遮蔽后代码可编译,但担心新代码中单独await请求的逻辑仍会导致无限循环,无法解决最初的问题。现寻求正确的超时实现方案,以及如何避免无限等待的问题。

原请求循环代码

fn construct_headers() -> HeaderMap {
    let mut headers = HeaderMap::new();
    headers.insert(USER_AGENT, HeaderValue::from_static("xxx-zzz"));
    headers
}

let mut headers = HeaderMap::new(); headers.insert(USER_AGENT, HeaderValue::from_static("reqwest"));
let client = reqwest::Client::builder().build()?;
let params = [("subreddit", subr), ("size", &size.to_string()), ("before", &last.clone())];

let mut resp = client.get("https://api.pushshift.io/reddit/search/submission/")
    .query(&params)
    .headers(construct_headers())
    .send()
    .await?
    .text()
    .await?
;

if resp.to_string().contains("Too Many"){
    'rqst:loop{
        thread::sleep(time::Duration::from_secs(1));
        resp = client.get("https://api.pushshift.io/reddit/search/submission/")
            .query(&params)
            .headers(construct_headers())
            .send()
            .await?
            .text()
            .await?
        ;
        if resp.to_string().contains("Too Many"){
            continue
        } else {
            break 'rqst
        };
    };
};

添加超时后的错误代码

if resp.to_string().contains("Too Many"){
    'rqst:loop{
        thread::sleep(time::Duration::from_secs(1));
        resp = client.get("https://api.pushshift.io/reddit/search/submission/")
            .query(&params)
            .headers(construct_headers())
            .send()
        ;
        match tokio::time::timeout(std::time::Duration::from_secs(30), resp.send()).await {
            Ok(result) => match result {
                Ok(response) => response.text().await?, 
                Err(e) => return Ok(linkz), 
            },
            Err(_) => return Ok(linkz), 
        };
    if resp.to_string().contains("Too Many"){
        continue
    } else {
        break 'rqst
    };
    };
};

编译错误信息

error[E0308]: mismatched types
   --> src/main.rs:204:24
    |
191 |           let mut resp = client.get("https://api.pushshift.io/reddit/search/submission/")
    |  ________________________-
192 | |             .query(&params)
193 | |             .headers(construct_headers())
194 | |             .send()
195 | |             .await?
196 | |             .text()
197 | |             .await?
    | |___________________- expected due to this value


204 |                   resp = client.get("https://api.pushshift.io/reddit/search/submission/")
    |  ________________________^
205 | |                     .query(&params)
206 | |                     .headers(construct_headers())
207 | |                     .send()
    | |___________________________^ expected struct `std::string::String`, found opaque type
...
error[E0277]: `std::string::String` is not a future
   --> src/main.rs:210:80
    |
210 |                 match tokio::time::timeout(std::time::Duration::from_secs(30), resp).await {
    |                       --------------------                                     ^^^^ `std::string::String` is not a future

修改后的代码

if resp.to_string().contains("Too Many"){
    println!("2many rqstz");
    'rqst:loop{
        thread::sleep(time::Duration::from_secs(1));
        let resp = client.get("https://api.pushshift.io/reddit/search/submission/")
            .query(&params)
            .headers(construct_headers())
            .send()
        ;
        match tokio::time::timeout(std::time::Duration::from_secs(30), async {&resp}).await {
            Ok(result) => result,
            Err(_) => return Ok(linkz),
        };
    if resp.await?.text().await?.contains("Too Many"){
        println!("2mny");
        continue
    } else {
        break 'rqst
    };
    };
};

错误原因分析

  1. 类型不匹配:初始的resp是String类型(通过.text().await?解析得到),但后续直接将send()返回的SendFuture赋值给它,两种类型完全不同,触发编译错误。
  2. 逻辑错误:错误代码中调用resp.send(),但此时resp已是String,根本不存在send方法。
  3. 无效超时包裹:修改后的代码中,timeout仅包裹了async {&resp}(Future的引用),实际请求的await逻辑在超时范围外,仍可能出现无限等待。

正确实现方案

核心思路是将完整的请求+解析逻辑封装为异步函数,用timeout包裹整个函数调用,同时替换阻塞的thread::sleep为异步的tokio::time::sleep,避免阻塞异步runtime。

完整修正代码

use reqwest::{Client, HeaderMap, HeaderValue};
use tokio::time::{sleep, timeout, Duration};

const USER_AGENT: &str = "xxx-zzz";

fn construct_headers() -> HeaderMap {
    let mut headers = HeaderMap::new();
    headers.insert(USER_AGENT, HeaderValue::from_static("xxx-zzz"));
    headers
}

// 封装完整的请求+解析逻辑为异步函数
async fn fetch_submissions(client: &Client, params: &[(&str, &str)]) -> Result<String, reqwest::Error> {
    client.get("https://api.pushshift.io/reddit/search/submission/")
        .query(params)
        .headers(construct_headers())
        .send()
        .await?
        .text()
        .await
}

// 主逻辑示例
async fn main_logic(subr: &str, size: u32, last: &str) -> Result<(), Box<dyn std::error::Error>> {
    let client = Client::builder().build()?;
    let params = [("subreddit", subr), ("size", &size.to_string()), ("before", last)];

    // 初始请求
    let mut resp = fetch_submissions(&client, &params).await?;

    // 处理Too Many请求的循环
    while resp.contains("Too Many") {
        println!("2many rqstz");
        // 异步睡眠,不阻塞runtime
        sleep(Duration::from_secs(1)).await;
        
        // 用timeout包裹整个请求逻辑,30秒超时
        match timeout(Duration::from_secs(30), fetch_submissions(&client, &params)).await {
            Ok(Ok(new_resp)) => {
                resp = new_resp;
                // 检查是否还触发Too Many
                if !resp.contains("Too Many") {
                    break;
                }
            }
            Ok(Err(req_err)) => {
                // 请求本身出错,可根据需求重试或返回
                eprintln!("请求错误: {}", req_err);
                continue;
            }
            Err(_timeout_err) => {
                // 超时处理,比如返回或重试
                eprintln!("请求超时");
                return Ok(()); // 可根据需求调整,比如改为继续重试
            }
        }
    }

    // 后续处理resp...
    Ok(())
}

关键改进点

  1. 封装请求逻辑:将get -> send -> text的完整流程封装为fetch_submissions函数,确保每次请求逻辑一致,也方便用timeout整体包裹。
  2. 正确使用超时:timeout直接包裹fetch_submissions的调用,确保整个请求+解析过程都受超时限制,避免无限等待。
  3. 异步睡眠:用tokio::time::sleep代替thread::sleep,避免阻塞异步线程,让其他任务正常运行。
  4. 类型一致性:resp始终保持String类型,每次请求成功后更新为新的响应文本,避免类型不匹配问题。

额外建议

  • 添加重试次数限制,避免无限循环重试。
  • 针对Pushshift API的限流,建议遵循其速率限制,比如调整睡眠时间或使用指数退避策略。
  • 可结合reqwest客户端自带的超时配置(Client::builder().timeout(Duration::from_secs(30))),但tokio::time::timeout能提供更上层的超时控制(比如包含重试等待的总时长)。

内容的提问来源于stack exchange,提问作者tommy1213

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.25 23:18:26