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

如何用tokio-tungstenite与graphql_client并行执行GraphQL订阅及HTTP请求

解决方案

核心问题在于你用tokio::block_on阻塞了线程,导致异步任务无法被调度;而tokio::join仅等待任务完成一次,不适配周期性执行的场景。正确的做法是利用Tokio的异步任务调度能力,将WebSocket订阅和定时突变分别作为独立异步任务并行运行。

实现步骤

  1. 依赖Tokio运行时宏
    用#[tokio::main]注解主函数,自动启动Tokio异步运行时,无需手动调用block_on阻塞线程。

  2. Spawn独立的WebSocket订阅任务
    将WebSocket连接、订阅初始化和消息循环包装成异步任务,通过tokio::spawn提交给运行时调度,避免阻塞主线程或其他任务。

  3. 用Interval实现定时突变
    利用tokio::time::interval创建周期性定时器,在循环中每隔一分钟执行一次突变请求,同样通过tokio::spawn作为独立任务运行。

完整代码示例

use tokio_tungstenite::{connect_async, tungstenite::Message};
use serde_json::json;
use reqwest::Client;

#[tokio::main]
async fn main() {
    // 启动WebSocket订阅任务
    let ws_subscription_task = tokio::spawn(async {
        // 连接GraphQL WebSocket端点
        let (ws_stream, _resp) = connect_async("wss://your-graphql-subscription-endpoint")
            .await
            .expect("Failed to connect to WebSocket endpoint");

        let (mut write, mut read) = ws_stream.split();

        // 发送订阅初始化消息
        let subscribe_msg = json!({
            "type": "start",
            "id": "sub_1",
            "payload": {
                "query": "subscription { yourSubscriptionField { id, data } }"
            }
        });
        write.send(Message::Text(subscribe_msg.to_string()))
            .await
            .expect("Failed to send subscription request");

        // 循环接收订阅消息
        while let Some(msg_result) = read.next().await {
            match msg_result {
                Ok(msg) => {
                    if let Message::Text(text) = msg {
                        println!("Received subscription update: {}", text);
                        // 处理订阅数据的业务逻辑
                    }
                }
                Err(e) => {
                    eprintln!("WebSocket error: {}", e);
                    break;
                }
            }
        }
    });

    // 启动定时突变任务(每分钟执行一次)
    let scheduled_mutation_task = tokio::spawn(async {
        let mut interval = tokio::time::interval(tokio::time::Duration::from_secs(60));
        // 取消下面注释可让突变首次立即执行,否则等待1分钟后首次执行
        // interval.tick().await;

        let http_client = Client::new();
        let graphql_endpoint = "https://your-graphql-http-endpoint";

        loop {
            interval.tick().await;

            // 构造突变请求
            let mutation_payload = json!({
                "query": "mutation YourMutation($input: YourInputType!) { yourMutation(input: $input) { success } }",
                "variables": {
                    "input": { "field1": "value1", "field2": 2 }
                }
            });

            // 发送HTTP突变请求
            match http_client.post(graphql_endpoint)
                .header("Content-Type", "application/json")
                .json(&mutation_payload)
                .send()
                .await
            {
                Ok(resp) => {
                    let resp_text = resp.text().await.unwrap_or_default();
                    println!("Mutation executed successfully: {}", resp_text);
                }
                Err(e) => {
                    eprintln!("Mutation failed: {}", e);
                }
            }
        }
    });

    // 等待两个任务完成(无限循环,除非出错退出)
    let (ws_result, mutation_result) = tokio::join!(ws_subscription_task, scheduled_mutation_task);

    // 处理任务退出的错误
    if let Err(e) = ws_result {
        eprintln!("WebSocket task failed: {}", e);
    }
    if let Err(e) = mutation_result {
        eprintln!("Mutation task failed: {}", e);
    }
}

关键注意事项

  • 避免阻塞调用:全程使用Tokio提供的异步API(如connect_async、interval.tick()),不要用同步阻塞方法,确保任务能被运行时高效调度。
  • 任务隔离:两个任务通过tokio::spawn完全隔离,WebSocket的消息循环不会影响定时任务的执行,反之亦然。
  • 错误处理:示例中用expect和简单匹配处理错误,实际项目中建议使用更健壮的错误处理(如anyhow/thiserror crate),避免程序因单次错误崩溃。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 14:33:19