pub struct LocalSet { /* private fields */ }展开描述
在同一线程上执行的若干任务的集合。
在某些情况下,需要运行一个或多个未实现 Send 的 future,因此在线程之间发送它们是不安全的。在这些情况下,可以使用本地任务集合来调度一个或多个 !Send future 在同一线程上一起运行。
例如,以下代码将无法编译:
use std::rc::Rc;
#[tokio::main]
async fn main() {
// `Rc` does not implement `Send`, and thus may not be sent between
// threads safely.
let nonsend_data = Rc::new("my nonsend data...");
let nonsend_data = nonsend_data.clone();
// Because the `async` block here moves `nonsend_data`, the future is `!Send`.
// Since `tokio::spawn` requires the spawned future to implement `Send`, this
// will not compile.
tokio::spawn(async move {
println!("{}", nonsend_data);
// ...
}).await.unwrap();
}§Use with run_until
要派生 !Send future,我们可以使用本地任务集合,将它们调度在调用 Runtime::block_on 的线程上。在本地任务集合内运行时,我们可以使用 task::spawn_local,它可以派生 !Send future。例如:
use std::rc::Rc;
use tokio::task;
let nonsend_data = Rc::new("my nonsend data...");
// Construct a local task set that can run `!Send` futures.
let local = task::LocalSet::new();
// Run the local task set.
local.run_until(async move {
let nonsend_data = nonsend_data.clone();
// `spawn_local` ensures that the future is spawned on the local
// task set.
task::spawn_local(async move {
println!("{}", nonsend_data);
// ...
}).await.unwrap();
}).await;注意:run_until 方法只能在 #[tokio::main]、#[tokio::test] 或直接调用 Runtime::block_on 内部使用。它不能在使用 tokio::spawn 派生的任务内部使用。
§Awaiting a LocalSet
此外,LocalSet 本身实现了 Future,当 LocalSet 上所有派生任务都完成时,它也会随之完成。这可用于在 LocalSet 上运行多个 future,并驱动整个集合直到它们全部完成。例如:
use tokio::{task, time};
use std::rc::Rc;
let nonsend_data = Rc::new("world");
let local = task::LocalSet::new();
let nonsend_data2 = nonsend_data.clone();
local.spawn_local(async move {
// ...
println!("hello {}", nonsend_data2)
});
local.spawn_local(async move {
time::sleep(time::Duration::from_millis(100)).await;
println!("goodbye {}", nonsend_data)
});
// ...
local.await;注意:对 LocalSet 的 await 只能在 #[tokio::main]、#[tokio::test] 或直接调用 Runtime::block_on 内部进行。它不能在使用 tokio::spawn 派生的任务内部使用。
§Use inside tokio::spawn
上面提到的两种方法都不能在 tokio::spawn 内部使用,因此要在 tokio::spawn 内部派生 !Send future,我们需要采用其他方式。解决方案是在其他地方创建 LocalSet,并通过 mpsc 通道与它通信。
以下示例将 LocalSet 放在一个新线程中。
use tokio::runtime::Builder;
use tokio::sync::{mpsc, oneshot};
use tokio::task::LocalSet;
// This struct describes the task you want to spawn. Here we include
// some simple examples. The oneshot channel allows sending a response
// to the spawner.
#[derive(Debug)]
enum Task {
PrintNumber(u32),
AddOne(u32, oneshot::Sender<u32>),
}
#[derive(Clone)]
struct LocalSpawner {
send: mpsc::UnboundedSender<Task>,
}
impl LocalSpawner {
pub fn new() -> Self {
let (send, mut recv) = mpsc::unbounded_channel();
let rt = Builder::new_current_thread()
.enable_all()
.build()
.unwrap();
std::thread::spawn(move || {
let local = LocalSet::new();
local.spawn_local(async move {
while let Some(new_task) = recv.recv().await {
tokio::task::spawn_local(run_task(new_task));
}
// If the while loop returns, then all the LocalSpawner
// objects have been dropped.
});
// This will return once all senders are dropped and all
// spawned tasks have returned.
rt.block_on(local);
});
Self {
send,
}
}
pub fn spawn(&self, task: Task) {
self.send.send(task).expect("Thread with LocalSet has shut down.");
}
}
// This task may do !Send stuff. We use printing a number as an example,
// but it could be anything.
//
// The Task struct is an enum to support spawning many different kinds
// of operations.
async fn run_task(task: Task) {
match task {
Task::PrintNumber(n) => {
println!("{}", n);
},
Task::AddOne(n, response) => {
// We ignore failures to send the response.
let _ = response.send(n + 1);
},
}
}
#[tokio::main]
async fn main() {
let spawner = LocalSpawner::new();
let (send, response) = oneshot::channel();
spawner.spawn(Task::AddOne(10, send));
let eleven = response.await.unwrap();
assert_eq!(eleven, 11);
}实现§
Source§impl LocalSet
impl LocalSet
Sourcepub fn enter(&self) -> LocalEnterGuard
pub fn enter(&self) -> LocalEnterGuard
进入此 LocalSet 的上下文。
spawn_local 方法会将任务派生到你所进入的 LocalSet 上下文中。
Sourcepub fn spawn_local<F>(&self, future: F) -> JoinHandle<F::Output> ⓘ
pub fn spawn_local<F>(&self, future: F) -> JoinHandle<F::Output> ⓘ
将一个 !Send 任务派生到本地任务集合上。
此任务保证在当前线程上运行。
与自由函数 spawn_local 不同,此方法可用于在 LocalSet 没有运行时派生本地任务。提供的 future 将在 LocalSet 下次启动时开始运行,即使你没有 await 返回的 JoinHandle。
§示例
use tokio::task;
let local = task::LocalSet::new();
// Spawn a future on the local set. This future will be run when
// we call `run_until` to drive the task set.
local.spawn_local(async {
// ...
});
// Run the local task set.
local.run_until(async move {
// ...
}).await;
// When `run` finishes, we can spawn _more_ futures, which will
// run in subsequent calls to `run_until`.
local.spawn_local(async {
// ...
});
local.run_until(async move {
// ...
}).await;Sourcepub fn block_on<F>(&self, rt: &Runtime, future: F) -> F::Outputwhere
F: Future,
pub fn block_on<F>(&self, rt: &Runtime, future: F) -> F::Outputwhere
F: Future,
在给定的 runtime 上运行 future 直到完成,同时在当前线程上驱动该任务集合上派生的所有本地 future。
此方法在 runtime 上运行给定的 future,阻塞直到其完成,并产出其已解析的结果。该 future 内部派生的任何任务或定时器都将在 runtime 上执行。该 future 也可以调用 spawn_local,在当前线程上派生其他本地 future。
不应在异步上下文中调用此方法。
§Panics
如果 executor 处于满载状态、如果提供的 future 发生 panic,或者在异步执行上下文中调用此函数,则此函数会 panic。
§Notes
由于此函数在内部调用 Runtime::block_on,并在 block_on 调用内部驱动本地任务集合中的 future,因此本地 future 不得使用原地阻塞。如果本地任务需要发起阻塞调用,可以使用 spawn_blocking API 来替代。
例如,这将发生 panic:
use tokio::runtime::Runtime;
use tokio::task;
let rt = Runtime::new().unwrap();
let local = task::LocalSet::new();
local.block_on(&rt, async {
let join = task::spawn_local(async {
let blocking_result = task::block_in_place(|| {
// ...
});
// ...
});
join.await.unwrap();
})然而,下面的情况不会发生 panic:
use tokio::runtime::Runtime;
use tokio::task;
let rt = Runtime::new().unwrap();
let local = task::LocalSet::new();
local.block_on(&rt, async {
let join = task::spawn_local(async {
let blocking_result = task::spawn_blocking(|| {
// ...
}).await;
// ...
});
join.await.unwrap();
})Sourcepub async fn run_until<F>(&self, future: F) -> F::Outputwhere
F: Future,
pub async fn run_until<F>(&self, future: F) -> F::Outputwhere
F: Future,
在本地集合上运行 future 直到完成,并返回其输出。
此方法返回一个 future,它会在本地集合下运行给定的 future,从而允许该 future 调用 spawn_local 来派生其他 !Send future。在传递给 run_until 的 future 完成之前,在本地集合上派生的任何本地 future 都将在后台被驱动。当传递给 run_until 的 future 结束时,任何尚未完成的本地 future 仍会保留在本地集合上,并将在后续对 run_until 的调用中或 await 本地集合本身时被驱动。
§Cancel safety
当 future 自身 cancel safe 时,此方法是 cancel safe 的。
§示例
use tokio::task;
task::LocalSet::new().run_until(async {
task::spawn_local(async move {
// ...
}).await.unwrap();
// ...
}).await;