【问题标题】:Async read from UdpSocket从 UdpSocket 异步读取
【发布时间】:2018-02-06 10:21:53
【问题描述】:

我正在尝试在 Tokio 中同时处理到达的 UDP 数据包。但是,以下 MWE 并没有达到我的预期:

extern crate futures;
extern crate tokio_core;
extern crate tokio_io;

use futures::{Future, Stream};
use std::net::SocketAddr;
use tokio_core::net::{UdpCodec, UdpSocket};
use tokio_core::reactor::Core;

// just a codec to send and receive bytes
pub struct LineCodec;
impl UdpCodec for LineCodec {
    type In = (SocketAddr, Vec<u8>);
    type Out = (SocketAddr, Vec<u8>);

    fn decode(&mut self, addr: &SocketAddr, buf: &[u8]) -> std::io::Result<Self::In> {
        Ok((*addr, buf.to_vec()))
    }

    fn encode(&mut self, (addr, buf): Self::Out, into: &mut Vec<u8>) -> SocketAddr {
        into.extend(buf);
        addr
    }
}

fn compute(addr: SocketAddr, msg: Vec<u8>) -> Box<Future<Item = (), Error = ()>> {
    println!("Starting to compute for: {}", addr);
    // sleep is a placeholder for a long computation
    std::thread::sleep(std::time::Duration::from_secs(8));
    println!("Done computing for for: {}", addr);
    Box::new(futures::future::ok(()))
}

fn main() {
    let mut core = Core::new().unwrap();
    let handle = core.handle();
    let listening_addr = "127.0.0.1:8080".parse::<SocketAddr>().unwrap();
    let socket = UdpSocket::bind(&listening_addr, &handle).unwrap();
    println!("Listening on: {}", socket.local_addr().unwrap());

    let (writer, reader) = socket.framed(LineCodec).split();

    let socket_read = reader.for_each(|(addr, msg)| {
        println!("Got {:?}", msg);
        handle.spawn(compute(addr, msg));
        Ok(())
    });

    core.run(socket_read).unwrap();
}

$ nc -u localhost 8080连接两个终端并发送一些文本,我可以看到来自第二个终端的消息是在第一个完成后处理的。

我需要改变什么?

【问题讨论】:

    标签: asynchronous concurrency rust future rust-tokio


    【解决方案1】:

    正如@Stefan 在另一个答案中所说,您不应该阻塞异步代码。鉴于您的示例,看起来 sleep 是一些长时间计算的占位符。因此,您应该将该计算委托给另一个线程,例如this example,而不是使用超时:

    extern crate futures;
    extern crate futures_cpupool;
    
    use futures::Future;
    use futures_cpupool::CpuPool;
    
    ...
    
    let pool = CpuPool::new_num_cpus();
    
    ...
    
    fn compute(handle: &Handle, addr: SocketAddr, _msg: Vec<u8>) -> Box<Future<Item = (), Error = ()>> {
        // I don't know enough about Tokio to know how to make `pool` available here
        pool.spawn_fn (|| {
            println!("Starting to compute for: {}", addr);
            std::thread::sleep(std::time::Duration::from_secs(8));
            println!("Done computing for for: {}", addr);
            Ok(())
        })
    }
    

    【讨论】:

    • 是的,这就是我的意思。我对我的问题发表了评论,以澄清睡眠的使用。
    【解决方案2】:

    永远不要在异步代码中使用sleep(并避免任何其他阻塞调用)。

    您可能希望像这样使用Timeout

    Playground

    fn compute(handle: &Handle, addr: SocketAddr, _msg: Vec<u8>) -> Box<Future<Item = (), Error = ()>> {
        println!("Starting to compute for: {}", addr);
        Box::new(
            Timeout::new(std::time::Duration::from_secs(8), handle)
                .unwrap()
                .map_err(|e| panic!("timeout failed: {:?}", e))
                .and_then(move |()| {
                    println!("Done computing for for: {}", addr);
                    Ok(())
                }),
        )
    }
    

    【讨论】:

      猜你喜欢
      • 2013-04-23
      • 2014-10-11
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2016-09-27
      • 2012-01-29
      • 2014-11-06
      • 2011-12-29
      相关资源
      最近更新 更多