跳到主内容
Rust
文章阅读

在Rust中使用异步futures

2025/12/1669 次阅读5 分钟

异步不是rust自带的标准库,所以要先添加依赖,后面要用异步控制流,把tokio也加上去 Cargo.toml


futures = "0.3.30"
tokio = { version = "1.35.1", features = ["full"] }
reqwest = "0.11.23"
anyhow = "1.0.79"

1.异步基础

基础语法,fn函数前面添加async,等待异步函数执行完毕,在函数末尾添加.await


use tokio::time;
async fn count_to(count: i32) {
    for i in 1..count {
        println!("count is : {i}!");
    }
}

async fn async_main(count: i32) {
    count_to(count).await
}

#[tokio::main]
async fn main() {
    tokio::spawn(count_tov2(10));
    for i in 1..5 {
        println!("main task: {i}");
        time::sleep(time::Duration::from_millis(5)).await;
    }
}

async fn count_tov2(count: i32) {
    for i in 1..count {
        println!("count in task : {i}?");
        time::sleep(time::Duration::from_millis(5)).await
    }
}

2.异步控制流

可以使用future::join_all来等待所有的任务完成,单个任务使用join


use anyhow::Result;
use futures::future;
use reqwest;
use std::collections::HashMap;
async fn size_of_page(url: &str) -> Result<usize> {
    let resp = reqwest::get(url).await?;
    Ok(resp.text().await?.len())
}

#[tokio::main]
async fn main() {
    let urls: [&str; 4] = [
        "https://baidu.com",
        "https://itbug.shop",
        "https://play.rust-lang.org/",
        "ABD_URL",
    ];
    let futures_iter = urls.into_iter().map(size_of_page);
    let results = future::join_all(futures_iter).await;
    let page_size_dict: HashMap<&str, Result<usize>> =
        urls.into_iter().zip(results.into_iter()).collect();
    println!("{:?}", page_size_dict)
}

3.异步通道

可以使用tokio::sync::mpsc通道来实现异步任务之间的数据传输和通讯


use tokio::sync::mpsc::{self, Receiver};

async fn ping_handle(mut input: Receiver<()>) {
    let mut count: usize = 0;
    while let Some(_) = input.recv().await {
        count += 1;
        println!("received {count} pings so far.");
    }
    println!("ping_handle complete");
}

#[tokio::main]
async fn main() {
    let (s, r) = mpsc::channel(32);

    let ping_handle_task = tokio::spawn(ping_handle(r));
    for i in 0..10 {
        s.send(()).await.expect("failed to send ping");
        println!("sent {} pings so far.", i + 1);
    }

    drop(s);
    ping_handle_task
        .await
        .expect("something went wrong in ping");
}

4. 异步任务

在异步任务中使用io流来对连接进行读写


use tokio::io::{self, AsyncReadExt, AsyncWriteExt};
use tokio::net::TcpListener;
#[tokio::main]
async fn main() -> io::Result<()> {
    let listener = TcpListener::bind("127.0.0.1:6142").await?;
    println!("listening on port 6142");
    loop {
        let (mut socket, addr) = listener.accept().await?;
        println!("connection form {addr:?}");
        tokio::spawn(async move {
            if let Err(e) = socket.write_all(b"who are you?\
").await {
                println!("socket error: {e:?}");
                return;
            }
            let mut bug = vec![0; 1024];
            let reply = match socket.read(&mut bug).await {
                Ok(n) => {
                    let name = std::str::from_utf8(&bug[..n]).unwrap().trim();
                    format!("thanks for dialing in,{name}!\
")
                }
                Err(e) => {
                    println!("socket error :{e:?}");
                    return;
                }
            };
            if let Err(e) = socket.write_all(reply.as_bytes()).await {
                println!("socket error {e:?}");
            }
        });
    }
}

返回顶部