异步不是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:?}");
}
});
}
}