【问题标题】:How can multiple threads share an iterator?多个线程如何共享一个迭代器?
【发布时间】:2017-07-25 19:32:33
【问题描述】:

我一直在研究一个函数,该函数将使用 Rust 和线程将一堆文件从源复制到目标。我在让线程共享迭代器时遇到了一些麻烦。借钱系统我还不习惯:

extern crate libc;
extern crate num_cpus;

use libc::{c_char, size_t};
use std::thread;
use std::fs::copy;

fn python_str_array_2_str_vec<T, U, V>(_: T, _: U) -> V {
    unimplemented!()
}

#[no_mangle]
pub extern "C" fn copyFiles(
    sources: *const *const c_char,
    destinies: *const *const c_char,
    array_len: size_t,
) {
    let src: Vec<&str> = python_str_array_2_str_vec(sources, array_len);
    let dst: Vec<&str> = python_str_array_2_str_vec(destinies, array_len);
    let mut iter = src.iter().zip(dst);
    let num_threads = num_cpus::get();
    let threads = (0..num_threads).map(|_| {
        thread::spawn(|| while let Some((s, d)) = iter.next() {
            copy(s, d);
        })
    });
    for t in threads {
        t.join();
    }
}

fn main() {}

我遇到了这个我无法解决的编译错误:

error[E0597]: `src` does not live long enough
  --> src/main.rs:20:20
   |
20 |     let mut iter = src.iter().zip(dst);
   |                    ^^^ does not live long enough
...
30 | }
   | - borrowed value only lives until here
   |
   = note: borrowed value must be valid for the static lifetime...

error[E0373]: closure may outlive the current function, but it borrows `**iter`, which is owned by the current function
  --> src/main.rs:23:23
   |
23 |         thread::spawn(|| while let Some((s, d)) = iter.next() {
   |                       ^^                          ---- `**iter` is borrowed here
   |                       |
   |                       may outlive borrowed value `**iter`
   |
help: to force the closure to take ownership of `**iter` (and any other referenced variables), use the `move` keyword, as shown:
   |         thread::spawn(move || while let Some((s, d)) = iter.next() {

我已经看过以下问题:

Value does not live long enough when using multiple threads 我没有使用chunks,我想尝试通过线程共享一个迭代器,尽管创建块将它们传递给线程将是经典的解决方案。

Unable to send a &str between threads because it does not live long enough 我已经看到了一些使用通道与线程通信的答案,但我不太确定使用它们。应该有一种更简单的方法来通过线程共享一个对象。

Why doesn't a local variable live long enough for thread::scoped 这引起了我的注意,scoped 应该修复我的错误,但由于它位于不稳定的通道中,我想看看是否有另一种方法可以使用 spawn。

有人可以解释我应该如何修复生命周期以便可以从线程中访问迭代器吗?

【问题讨论】:

  • @Shepmaster,所有线程都将共享耻辱迭代器,对于块解决方案,所有线程都有不同的迭代器。但是,如果是最便宜和/或最简单的解决方案,我会继续这样做:)
  • @DanielSanchez,您可以使用 Arc&lt;Mutex&lt;...&gt;&gt; Playground 共享迭代器
  • 您的代码实际上不会并行运行:它会创建并立即一个接一个地加入线程,因为map 是惰性的。你应该使用 for 循环,或者collect 结果。
  • @interjay,好的,我将尝试收集并仅使用块并对其进行迭代以向线程发送垃圾邮件。谢谢!

标签: multithreading rust borrowing


【解决方案1】:

这是您的问题的minimal, reproducible example:

use std::thread;

fn main() {
    let src = vec!["one"];
    let dst = vec!["two"];
    let mut iter = src.iter().zip(dst);
    thread::spawn(|| {
        while let Some((s, d)) = iter.next() {
            println!("{} -> {}", s, d);
        }
    });
}

有多个相关问题:

  1. 迭代器位于堆栈中,线程的闭包会引用它。
  2. 闭包采用对迭代器的可变引用。
  3. 迭代器本身具有对位于堆栈中的Vec 的引用。
  4. Vec 本身有对字符串切片的引用,这些切片可能存在于堆栈中,但不能保证其存在时间比线程长。

换句话说,Rust 编译器阻止了你执行四个独立的内存不安全部分。

要认识到的主要一点是,您生成的任何线程可能比您生成它的地方寿命长。即使您立即调用join,编译器也无法静态验证会发生这种情况,因此它必须采取保守的路径。这就是scoped threads 的重点——它们保证线程在它们开始的堆栈帧之前退出。

此外,您正尝试在多个并发线程中使用可变引用。 零保证可以安全地并行调用迭代器(或构建它的任何迭代器)。两个线程完全有可能恰好同时调用next。这两段代码并行运行并写入相同的内存地址。一个线程写入一半数据,另一个线程写入另一半,现在您的程序在未来的某个任意时间点崩溃。

使用crossbeam 之类的工具,您的代码将如下所示:

use crossbeam; // 0.7.3

fn main() {
    let src = vec!["one"];
    let dst = vec!["two"];

    let mut iter = src.iter().zip(dst);
    while let Some((s, d)) = iter.next() {
        crossbeam::scope(|scope| {
            scope.spawn(|_| {
                println!("{} -> {}", s, d);
            });
        })
        .unwrap();
    }
}

如前所述,这一次只会产生一个线程,等待它完成。获得更多并行性的另一种方法(本练习的通常要点)是交换对next 和spawn 的调用。这需要通过move 关键字将s 和d 的所有权转移给线程:

use crossbeam; // 0.7.3

fn main() {
    let src = vec!["one", "alpha"];
    let dst = vec!["two", "beta"];

    let mut iter = src.iter().zip(dst);
    crossbeam::scope(|scope| {
        while let Some((s, d)) = iter.next() {
            scope.spawn(move |_| {
                println!("{} -> {}", s, d);
            });
        }
    })
    .unwrap();
}

如果您在spawn 中添加睡眠调用,您可以看到线程并行运行。

我会使用 for 循环来编写它,但是:

let iter = src.iter().zip(dst);
crossbeam::scope(|scope| {
    for (s, d) in iter {
        scope.spawn(move |_| {
            println!("{} -> {}", s, d);
        });
    }
}).unwrap();

最后,迭代器在当前线程上运行,然后从迭代器返回的每个值都被移交给一个新线程。保证新线程在捕获的引用之前退出。

您可能对Rayon 感兴趣,这是一个允许某些类型的迭代器轻松并行化的板条箱。

另见:

【讨论】:

  • 非常感谢,它真的让我对 Rust 有了更多的了解。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2011-10-20
  • 1970-01-01
  • 2014-03-08
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2015-05-08
相关资源
最近更新 更多