【问题标题】:Why do I not get a wakeup for multiple futures when they use the same underlying socket?当多个期货使用相同的底层套接字时,为什么我没有唤醒它们?
【发布时间】:2020-04-03 11:33:55
【问题描述】:

我的代码使用相同的本地 UdpSocket 将数据发送到多个 UDP 端点:

use futures::stream::FuturesUnordered;
use futures::StreamExt;
use std::{
    future::Future,
    net::{Ipv4Addr, SocketAddr},
    pin::Pin,
    task::{Context, Poll},
};
use tokio::net::UdpSocket;

#[tokio::main]
async fn main() {
    let server_0: SocketAddr = (Ipv4Addr::UNSPECIFIED, 12000).into();
    let server_2: SocketAddr = (Ipv4Addr::UNSPECIFIED, 12002).into();
    let server_1: SocketAddr = (Ipv4Addr::UNSPECIFIED, 12001).into();

    tokio::spawn(start_server(server_0));
    tokio::spawn(start_server(server_1));
    tokio::spawn(start_server(server_2));

    let client_addr: SocketAddr = (Ipv4Addr::UNSPECIFIED, 12004).into();
    let socket = UdpSocket::bind(client_addr).await.unwrap();

    let mut futs = FuturesUnordered::new();
    futs.push(Task::new(0, &socket, &server_0));
    futs.push(Task::new(1, &socket, &server_1));
    futs.push(Task::new(2, &socket, &server_2));

    while let Some(n) = futs.next().await {
        println!("Done: {:?}", n)
    }
}

async fn start_server(addr: SocketAddr) {
    let mut socket = UdpSocket::bind(addr).await.unwrap();
    let mut buf = [0; 512];
    loop {
        println!("{:?}", socket.recv_from(&mut buf).await);
    }
}

struct Task<'a> {
    value: u32,
    socket: &'a UdpSocket,
    addr: &'a SocketAddr,
}

impl<'a> Task<'a> {
    fn new(value: u32, socket: &'a UdpSocket, addr: &'a SocketAddr) -> Self {
        Self {
            value,
            socket,
            addr,
        }
    }
}

impl Future for Task<'_> {
    type Output = Option<u32>;

    fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
        println!("Polling for {}", self.value);
        let buf = &self.value.to_be_bytes();

        match self.socket.poll_send_to(cx, buf, self.addr) {
            Poll::Ready(Ok(_)) => {
                println!("Got Ok for {}", self.value);
                Poll::Ready(Some(self.value))
            }
            Poll::Ready(Err(_)) => {
                println!("Got err for {}", self.value);
                Poll::Ready(None)
            }
            Poll::Pending => {
                println!("Got pending for {}", self.value);
                Poll::Pending
            }
        }
    }
}

有时只写一个数据就卡住了,打印:

Polling for 0
Got pending for 0
Polling for 1
Got pending for 1
Polling for 2
Got pending for 2
Polling for 2
Got Ok for 2
Done: Some(2)
Ok((4, V4(127.0.0.1:12004)))

在这种情况下,值 0 和 1 的任务永远不会被唤醒。我如何可靠地向他们发出信号以唤醒他们?

我尝试在收到Poll::Ready 时致电cx.waker().wake_by_ref(),因为我认为这也可能唤醒其他人,但事实并非如此。

【问题讨论】:

  • 事实证明你是对的,有时它也不适用于我的解决方案,好的,那么你有没有试过cx.waker().wake_by_ref()在接收Poll::Pending时调用它,因为Waker只会唤醒当前任务,而不是其他人。如果您的未来等待很多,这会带来高 CPU 使用率,这是一种解决方法。作为一个简单的解决方案,您可以在具有超时的单独线程上调用这些唤醒器(我的意思是轮询)。 Reference...&Waker,可用于唤醒当前任务。

标签: rust future rust-tokio


【解决方案1】:

poll_send_to 返回Poll::Pending 时,它保证向上下文中提供的Waker 发出唤醒信号。然而,它只需要发出一个唤醒到它被轮询的最后一个Waker。这意味着,由于您在多个任务的同一个套接字上调用poll_send_to,因此该套接字只承诺向最后轮询它的那个发出唤醒。

这也解释了为什么会这样:

let mut futs = Vec::new();
futs.push(Task::new(0, &socket, &server_0));
futs.push(Task::new(1, &socket, &server_1));
futs.push(Task::new(2, &socket, &server_2));

for n in join_all(futs).await {
    println!("Done: {:?}", n)
}

FuturesUnordered 不同,join_all 组合器每次轮询时都会轮询每个内部未来,但FuturesUnordered 会跟踪唤醒来自哪个潜在未来。

另见this thread

【讨论】:

  • 但这并不提供相同的行为,对吧?它将等待所有期货完成以产生结果?
  • 是的,它只是为了帮助解释问题所在而添加的。有关替代方案的讨论可以在链接的线程中进行。
  • 当然我读完了,谢谢,但我有一个问题,谁在 tokio executor 的第一次投票后唤醒JoinAll 任务?确实,内部任务不这样做是因为我们有这个问题?
  • 仍然会有三个任务之一收到通知。问题是JoinAll 将在通知其中任何一个任务时轮询内部的所有任务,因此与FuturesUnordered 不同的是,只通知其中一个任务是没有意义的;他们仍然会被轮询。
  • 那么这意味着如果所有期货在第一次投票时都返回待处理,它就不会起作用?如本例所示 ? play.rust-lang.org/…
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2011-06-13
  • 1970-01-01
  • 2019-06-07
  • 2012-01-25
  • 2013-02-08
  • 2021-12-28
  • 1970-01-01
相关资源
最近更新 更多