【问题标题】:How to create a ring communication between threads using mpsc channels?如何使用 mpsc 通道在线程之间创建环形通信?
【发布时间】:2020-09-26 16:47:22
【问题描述】:

我想生成 n 个线程,这些线程能够与环形拓扑中的其他线程进行通信,例如线程 0 可以向线程 1 发送消息,线程 1 向线程 2 发送消息,以此类推,线程 n 向线程 0 发送消息。

这是我想用 n=3 实现的示例:

use std::sync::mpsc::{self, Receiver, Sender};
use std::thread;

let (tx0, rx0): (Sender<i32>, Receiver<i32>) = mpsc::channel();
let (tx1, rx1): (Sender<i32>, Receiver<i32>) = mpsc::channel();
let (tx2, rx2): (Sender<i32>, Receiver<i32>) = mpsc::channel();

let child0 = thread::spawn(move || {
    tx0.send(0).unwrap();
    println!("thread 0 sent: 0");
    println!("thread 0 recv: {:?}", rx2.recv().unwrap());
});
let child1 = thread::spawn(move || {
    tx1.send(1).unwrap();
    println!("thread 1 sent: 1");
    println!("thread 1 recv: {:?}", rx0.recv().unwrap());
});
let child2 = thread::spawn(move || {
    tx2.send(2).unwrap();
    println!("thread 2 sent: 2");
    println!("thread 2 recv: {:?}", rx1.recv().unwrap());
});

child0.join();
child1.join();
child2.join();

在这里,我在循环中创建通道,将它们存储在向量中,重新排序发送者,将它们存储在新向量中,然后生成线程,每个线程都有自己的发送器-接收器(tx1/rx0、tx2/rx1 等)对。

const NTHREADS: usize = 8;

// create n channels
let channels: Vec<(Sender<i32>, Receiver<i32>)> =
    (0..NTHREADS).into_iter().map(|_| mpsc::channel()).collect();

// switch tupel entries for the senders to create ring topology
let mut channels_ring: Vec<(Sender<i32>, Receiver<i32>)> = (0..NTHREADS)
    .into_iter()
    .map(|i| {
        (
            channels[if i < channels.len() - 1 { i + 1 } else { 0 }].0,
            channels[i].1,
        )
    })
    .collect();

let mut children = Vec::new();
for i in 0..NTHREADS {
    let (tx, rx) = channels_ring.remove(i);

    let child = thread::spawn(move || {
        tx.send(i).unwrap();
        println!("thread {} sent: {}", i, i);
        println!("thread {} recv: {:?}", i, rx.recv().unwrap());
    });

    children.push(child);
}

for child in children {
    let _ = child.join();
}

这不起作用,因为无法复制 Sender 以创建新向量。 但是,如果我使用 refs (& Sender):

let mut channels_ring: Vec<(&Sender<i32>, Receiver<i32>)> = (0..NTHREADS)
    .into_iter()
    .map(|i| {
        (
            &channels[if i < channels.len() - 1 { i + 1 } else { 0 }].0,
            channels[i].1,
        )
    })
    .collect();

我无法生成线程,因为std::sync::mpsc::Sender&lt;i32&gt; 无法在线程之间安全地共享。

【问题讨论】:

    标签: multithreading rust channel


    【解决方案1】:

    Senders 和 Receivers 无法共享,因此您需要将它们移动到各自的线程中。这意味着将它们从Vec 中删除,或者在迭代它时消耗Vec - 即使作为中间步骤,向量也不允许处于无效状态(有孔)。使用into_iter 遍历向量将通过使用它们来实现。

    你可以用来让发送者和接收者在一个循环中配对的一个小技巧是创建两个向量;发送者之一和接收者之一;然后旋转一个,以便每个向量中的相同索引将为您提供所需的对。

    use std::sync::mpsc::{self, Receiver, Sender};
    use std::thread;
    
    fn main() {
        const NTHREADS: usize = 8;
    
        // create n channels
        let (mut senders, receivers): (Vec<Sender<i32>>, Vec<Receiver<i32>>) =
            (0..NTHREADS).into_iter().map(|_| mpsc::channel()).unzip();
    
        // move the first sender to the back
        senders.rotate_left(1);
    
        let children: Vec<_> = senders
            .into_iter()
            .zip(receivers.into_iter())
            .enumerate()
            .map(|(i, (tx, rx))| {
                thread::spawn(move || {
                    tx.send(i as i32).unwrap();
                    println!("thread {} sent: {}", i, i);
                    println!("thread {} recv: {:?}", i, rx.recv().unwrap());
                })
            })
            .collect();
    
        for child in children {
            let _ = child.join();
        }
    }
    

    【讨论】:

    • 美丽的答案。您可能还想使用collect 构造children,正如我在回答中所做的那样。
    • 是的,这样会更好。不过我会保留它,而不是对 OP 的原始代码进行太多更改。
    • 好吧,实际上它现在会困扰我,所以我必须改变它;)
    • 我认为这是有道理的,因为您(正确地)消除了大部分 OP 的代码,但这取决于您。 :)
    • 我喜欢你使用两个向量的巧妙技巧,然后旋转一个,然后将它们拉回一起。它避免了我最初将条目从一个向量复制到另一个向量的问题。但是,您的答案是以与我的意图相反的顺序链接线程。将 rotate_right 更改为 rotate_left 可以解决此问题。 (我也认为您忘记在编辑中省略子变量)
    【解决方案2】:

    这不起作用,因为无法复制 Sender 以创建新向量。但是,如果我使用 refs (& Sender):

    虽然Sender 确实无法复制,但它确实实现了Clone,因此您始终可以手动克隆它。但这种方法不适用于Receiver,它不是Clone,您还需要从向量中提取它。

    您的第一个代码的问题是您不能使用let foo = vec[i] 将一个值移出非Copy 值的向量。这将使向量处于无效状态,其中一个元素无效,随后对其进行访问将导致未定义的行为。为此,Vec 需要跟踪哪些元素被移动,哪些没有,这将对所有Vecs 施加成本。因此,Vec 不允许将元素移出其中,而是留给用户跟踪移动。

    将值移出Vec 的一种简单方法是将Vec&lt;T&gt; 替换为Vec&lt;Option&lt;T&gt;&gt; 并使用Option::takefoo = vec[i] 被替换为 foo = vec[i].take().unwrap(),这会将 T 值从 vec[i] 的选项中移出(同时断言它不是 None),并在向量。这是您以这种方式修改的第一次尝试 (playground):

    const NTHREADS: usize = 8;
    
    let channels_ring: Vec<_> = {
        let mut channels: Vec<_> = (0..NTHREADS)
            .into_iter()
            .map(|_| {
                let (tx, rx) = mpsc::channel();
                (Some(tx), Some(rx))
            })
            .collect();
    
        (0..NTHREADS)
            .into_iter()
            .map(|rxpos| {
                let txpos = if rxpos < NTHREADS - 1 { rxpos + 1 } else { 0 };
                (
                    channels[txpos].0.take().unwrap(),
                    channels[rxpos].1.take().unwrap(),
                )
            })
            .collect()
    };
    
    let children: Vec<_> = channels_ring
        .into_iter()
        .enumerate()
        .map(|(i, (tx, rx))| {
            thread::spawn(move || {
                tx.send(i as i32).unwrap();
                println!("thread {} sent: {}", i, i);
                println!("thread {} recv: {:?}", i, rx.recv().unwrap());
            })
        })
        .collect();
    
    for child in children {
        child.join().unwrap();
    }
    

    【讨论】:

    • 我喜欢这个答案,因为使用智能指针比彼得的答案更接近我的初始帖子。然而,我更喜欢彼得的优雅方法。虽然你先回答,但我接受彼得的回答可以吗?
    • @jhscheer 当然,选择接受哪个答案完全取决于您作为问题的作者,无论它们到达的顺序如何。我留下了我的答案,因为它实际上提供了一个不同的解决方案,这可能对在彼得的方法未涵盖的场景中解决相同问题的未来读者有用。
    • @jhscheer 另请注意,我的答案中没有智能指针,只有Option 便宜得多,因为它不会引入新的分配或间接。在this case 中,它甚至不会使向量变大,因为Rust optimizes enums for space 尽可能地。
    • @jhscheer 支持和接受答案不是一回事。如果它对您有帮助,您也可以投票赞成这个答案。
    • 再想一想,即使我最终使用了彼得的想法,我也会接受这个答案。因为它 i) 更接近我最初的方法,并且 ii) 帮助我更好地理解我的初始代码失败的原因以及将来如何解决类似问题。
    猜你喜欢
    • 1970-01-01
    • 2021-07-09
    • 2019-01-30
    • 1970-01-01
    • 1970-01-01
    • 2019-12-11
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多