【问题标题】:How to start and stop a worker thread如何启动和停止工作线程
【发布时间】:2021-01-03 08:23:31
【问题描述】:

我有以下要求,这在其他编程语言中是标准的,但我不知道如何在 Rust 中做。

我有一个类,我想编写一个方法来产生一个满足 2 个条件的工作线程:

  • 产生工作线程后,函数返回(所以其他地方不需要等待)
  • 有一种机制可以停止此线程。

例如,这是我的虚拟代码:

struct A {
    thread: JoinHandle<?>,
}

impl A {
    pub fn run(&mut self) -> Result<()>{
        self.thread = thread::spawn(move || {
            let mut i = 0;
            loop {
                self.call();
                i = 1 + i;
                if i > 5 {
                    return
                }
            }
        });
        Ok(())
    }

    pub fn stop(&mut self) -> std::thread::Result<_> {
        self.thread.join()
    }

    pub fn call(&mut self) {
        println!("hello world");
    }
}

fn main() {
    let mut a = A{};
    a.run();
}

thread: JoinHandle&lt;?&gt; 出现错误。在这种情况下,线程的类型是什么。我的代码在启动和停止工作线程方面是否正确?

【问题讨论】:

  • 您能否详细说明fn call 在实际代码中的作用?因为您需要使用Mutex,这意味着如果您调用stop(),您的工作线程可能会死锁。
  • fn call中,我将访问结构A的一些属性,这就是为什么我需要在这里使用move。这够了吗?因为逻辑很复杂,能不能帮我详细说明一下这里还有什么需要说的?
  • 结构 A 有一些与 mpsc 包相关的属性以及一些类型为 Arc&lt;Mutex 的属性。我不知道这些是否会影响完成工作的方式?
  • 差不多,只是稍微澄清一下。 A 是否需要共享,或者共享 struct AData 是否足够?此外,在工作线程完成之前,您真的需要访问 attributes 吗?如果没有,那么fn stop(...) -&gt; AData(简体)就足够了吗?
  • 我想我理解潜在的问题。我会写几个例子作为答案。如果它们还不够,我会在之后澄清问题:)

标签: multithreading rust


【解决方案1】:

简而言之,join() 中的 T JoinHandle&lt;T&gt; 返回传递给 thread::spawn() 的闭包结果。 所以在你的情况下JoinHandle&lt;?&gt; 需要是JoinHandle&lt;()&gt;,因为你的闭包返回nothing,即() (unit)

除此之外,您的虚拟代码还包含一些其他问题。

  • run() 的返回类型不正确,至少需要为 Result&lt;(), ()&gt;
  • thread 字段必须是 Option&lt;JoinHandle&lt;()&gt; 才能处理fn stop(&amp;mut self),因为join() 使用JoinHandle
  • 但是,您尝试将 &amp;mut self 传递给闭包,这会带来更多问题,归结为多个可变引用
    • 这可以通过例如解决Mutex&lt;A&gt;。但是,如果您调用 stop(),则可能会导致死锁。

但是,由于它是虚拟代码,并且您在 cmets 中进行了澄清。让我试着用几个例子来澄清你的意思。 这包括我重写你的虚拟代码。

工人完成后的结果

如果您在工作线程运行时不需要访问数据,那么您可以创建一个新的struct WorkerData。然后在run() 中,您从A 复制/克隆您需要的数据(或者我已将其重命名为Worker)。然后在闭包中你终于再次返回data,所以你可以通过join()获取它。

use std::thread::{self, JoinHandle};

struct WorkerData {
    ...
}

impl WorkerData {
    pub fn call(&mut self) {
        println!("hello world");
    }
}

struct Worker {
    thread: Option<JoinHandle<WorkerData>>,
}

impl Worker {
    pub fn new() -> Self {
        Self { thread: None }
    }

    pub fn run(&mut self) {
        // Create `WorkerData` and copy/clone whatever is needed from `self`
        let mut data = WorkerData {};

        self.thread = Some(thread::spawn(move || {
            let mut i = 0;
            loop {
                data.call();
                i = 1 + i;
                if i > 5 {
                    // Return `data` so we get in through `join()`
                    return data;
                }
            }
        }));
    }

    pub fn stop(&mut self) -> Option<thread::Result<WorkerData>> {
        if let Some(handle) = self.thread.take() {
            Some(handle.join())
        } else {
            None
        }
    }
}

您实际上并不需要将thread 变为Option&lt;JoinHandle&lt;WorkerData&gt;&gt;,而只需使用JoinHandle&lt;WorkerData&gt;&gt;。因为如果您想再次调用run(),重新分配包含Worker 的变量会更容易。

所以现在我们可以简化Worker,删除Option,将stop改为使用thread,同时创建new() -&gt; Self代替run(&amp;mut self)

use std::thread::{self, JoinHandle};

struct Worker {
    thread: JoinHandle<WorkerData>,
}

impl Worker {
    pub fn new() -> Self {
        // Create `WorkerData` and copy/clone whatever is needed from `self`
        let mut data = WorkerData {};

        let thread = thread::spawn(move || {
            let mut i = 0;
            loop {
                data.call();
                i = 1 + i;
                if i > 5 {
                    return data;
                }
            }
        });

        Self { thread }
    }

    pub fn stop(self) -> thread::Result<WorkerData> {
        self.thread.join()
    }
}

分享WorkerData

如果您想在多个线程之间保留对WorkerData 的引用,那么您需要使用Arc。由于您还希望能够对其进行变异,因此您需要使用Mutex

如果您只在单个线程中进行变异,那么您也可以选择RwLock,与Mutex 相比,您可以同时锁定和获取多个不可变引用。

use std::sync::{Arc, RwLock};
use std::thread::{self, JoinHandle};

struct Worker {
    thread: JoinHandle<()>,
    data: Arc<RwLock<WorkerData>>,
}

impl Worker {
    pub fn new() -> Self {
        // Create `WorkerData` and copy/clone whatever is needed from `self`
        let data = Arc::new(RwLock::new(WorkerData {}));

        let thread = thread::spawn({
            let data = data.clone();
            move || {
                let mut i = 0;
                loop {
                    if let Ok(mut data) = data.write() {
                        data.call();
                    }

                    i = 1 + i;
                    if i > 5 {
                        return;
                    }
                }
            }
        });

        Self { thread, data }
    }

    pub fn stop(self) -> thread::Result<Arc<RwLock<WorkerData>>> {
        self.thread.join()?;
        // You might be able to unwrap and get the inner `WorkerData` here
        Ok(self.data)
    }
}

如果添加一个方法可以获取data,形式为Arc&lt;RwLock&lt;WorkerData&gt;&gt;。然后,如果您在调用stop() 之前克隆Arc 并锁定它(内部RwLock),那么这将导致死锁。 为避免这种情况,任何data() 方法都应返回&amp;WorkerData&amp;mut WorkerData 而不是Arc。这样你就无法调用stop() 并导致死锁。

标记停止工作人员

如果你真的想停止工作线程,那么你必须使用一个标志来通知它这样做。您可以以共享AtomicBool 的形式创建标志。

use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, RwLock};
use std::thread::{self, JoinHandle};

struct Worker {
    thread: JoinHandle<()>,
    data: Arc<RwLock<WorkerData>>,
    stop_flag: Arc<AtomicBool>,
}

impl Worker {
    pub fn new() -> Self {
        // Create `WorkerData` and copy/clone whatever is needed from `self`
        let data = Arc::new(RwLock::new(WorkerData {}));

        let stop_flag = Arc::new(AtomicBool::new(false));

        let thread = thread::spawn({
            let data = data.clone();
            let stop_flag = stop_flag.clone();
            move || {
                // let mut i = 0;
                loop {
                    if stop_flag.load(Ordering::Relaxed) {
                        break;
                    }

                    if let Ok(mut data) = data.write() {
                        data.call();
                    }

                    // i = 1 + i;
                    // if i > 5 {
                    //     return;
                    // }
                }
            }
        });

        Self {
            thread,
            data,
            stop_flag,
        }
    }

    pub fn stop(self) -> thread::Result<Arc<RwLock<WorkerData>>> {
        self.stop_flag.store(true, Ordering::Relaxed);
        self.thread.join()?;
        // You might be able to unwrap and get the inner `WorkerData` here
        Ok(self.data)
    }
}

多线程多任务

如果您想要处理多种任务,分布在多个线程中,那么这里有一个更通用的示例。

您已经提到使用mpsc。因此,您可以使用 SenderReceiver 以及自定义 TaskTaskResult 枚举。

use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::mpsc::{self, Receiver, Sender};
use std::sync::{Arc, Mutex};
use std::thread::{self, JoinHandle};

pub enum Task {
    ...
}

pub enum TaskResult {
    ...
}

pub type TaskSender = Sender<Task>;
pub type TaskReceiver = Receiver<Task>;

pub type ResultSender = Sender<TaskResult>;
pub type ResultReceiver = Receiver<TaskResult>;

struct Worker {
    threads: Vec<JoinHandle<()>>,
    task_sender: TaskSender,
    result_receiver: ResultReceiver,
    stop_flag: Arc<AtomicBool>,
}

impl Worker {
    pub fn new(num_threads: usize) -> Self {
        let (task_sender, task_receiver) = mpsc::channel();
        let (result_sender, result_receiver) = mpsc::channel();

        let task_receiver = Arc::new(Mutex::new(task_receiver));

        let stop_flag = Arc::new(AtomicBool::new(false));

        Self {
            threads: (0..num_threads)
                .map(|_| {
                    let task_receiver = task_receiver.clone();
                    let result_sender = result_sender.clone();
                    let stop_flag = stop_flag.clone();

                    thread::spawn(move || loop {
                        if stop_flag.load(Ordering::Relaxed) {
                            break;
                        }

                        let task_receiver = task_receiver.lock().unwrap();

                        if let Ok(task) = task_receiver.recv() {
                            drop(task_receiver);

                            // Perform the `task` here

                            // If the `Task` results in a `TaskResult` then create it and send it back
                            let result: TaskResult = ...;
                            // The `SendError` can be ignored as it only occurs if the receiver
                            // has already been deallocated
                            let _ = result_sender.send(result);
                        } else {
                            break;
                        }
                    })
                })
                .collect(),
            task_sender,
            result_receiver,
            stop_flag,
        }
    }

    pub fn stop(self) -> Vec<thread::Result<()>> {
        drop(self.task_sender);

        self.stop_flag.store(true, Ordering::Relaxed);

        self.threads
            .into_iter()
            .map(|t| t.join())
            .collect::<Vec<_>>()
    }

    #[inline]
    pub fn request(&mut self, task: Task) {
        self.task_sender.send(task).unwrap();
    }

    #[inline]
    pub fn result_receiver(&mut self) -> &ResultReceiver {
        &self.result_receiver
    }
}

使用Worker 以及发送任务和接收任务结果的示例如下所示:

fn main() {
    let mut worker = Worker::new(4);

    // Request that a `Task` is performed
    worker.request(task);

    // Receive a `TaskResult` if any are pending
    if let Ok(result) = worker.result_receiver().try_recv() {
        // Process the `TaskResult`
    }
}

在少数情况下,您可能需要为Task 和/或TaskResult 实现Send查看"Understanding the Send trait"

unsafe impl Send for Task {}
unsafe impl Send for TaskResult {}

【讨论】:

  • 这花费了相当多的时间,但我希望其中一些示例能为您和其他人提供帮助:)
  • 非常感谢。我希望我现在可以投票 x1000 次。我正在阅读,它真的对我有很大帮助。
  • 当您消化完答案后,如果有任何问题,请随时发表评论,或者如果它解决了您的问题/问题,请接受答案:)
  • 嗨。对于迟到的回复,我深表歉意。你的回答对我帮助很大。我有另一个问题,基于这个问题,但更明确,你能帮我回答吗? stackoverflow.com/questions/65654251/…
【解决方案2】:

JoinHandle 的类型参数应该是线程函数的返回类型。

在这种情况下,返回类型是一个空元组(),发音为unit。当只有一个可能的值时使用它,当没有指定返回类型时,它是函数的隐式“返回类型”。

你可以写JoinHandle&lt;()&gt;来表示该函数不会返回任何东西。

(注意:您的代码会遇到self.call() 的一些借用检查器问题,这可能需要使用Arc&lt;Mutex&lt;Self&gt;&gt; 来解决,但这是另一个问题。)

【讨论】:

  • run() 中的Result&lt;()&gt; 也不正确,同时thread 必须是Option&lt;JoinHandle&lt;()&gt;&gt; 才能处理stop :)
  • 我认为接下来会发生什么,因为self 已从工作线程中移除。你也能帮我带来一些启示吗?
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2021-01-14
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多