Структура Local Set
pub struct LocalSet { /* private fields */ }
rt.Набор задач, выполняемых в одном потоке.
В некоторых случаях необходимо запустить один или несколько фьючерсов, которые не реализуют Send, а значит, их небезопасно передавать между потоками. В таких случаях для запуска одного или нескольких фьючерсов !Send в одном потоке можно использовать набор локальных задач.
Например, следующий код не скомпилируется:
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();
}
Использование с run_until
Чтобы запускать фьючерсы !Send, мы можем использовать набор локальных задач, чтобы запланировать их выполнение в потоке, вызывающем Runtime::block_on. При работе внутри набора локальных задач можно использовать task::spawn_local, которая может запускать фьючерсы !Send. Например:
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.
Ожидание LocalSet
Кроме того, сам LocalSet реализует Future и завершается, когда завершаются все задачи, запущенные в LocalSet. Это можно использовать для запуска нескольких фьючерсов в LocalSet и ожидания завершения всего набора. Например:
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 можно только внутри #[tokio::main], #[tokio::test] или непосредственно внутри вызова Runtime::block_on. Это нельзя делать внутри задачи, запущенной с помощью tokio::spawn.
Использование внутри tokio::spawn
Два упомянутых выше метода нельзя использовать внутри tokio::spawn, поэтому для запуска фьючерсов !Send изнутри tokio::spawn нужно поступить иначе. Решение — создать 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);
}Реализации
impl LocalSet
pub fn enter(&self) -> LocalEnterGuard
Входит в контекст этого LocalSet.
Метод spawn_local запускает задачи в LocalSet, в контексте которого вы находитесь.
pub fn spawn_local<F>(&self, future: F) -> JoinHandle<F::Output> ⓘ
Запускает задачу !Send в наборе локальных задач.
Гарантируется, что эта задача будет выполняться в текущем потоке.
В отличие от свободной функции spawn_local, этот метод можно использовать для запуска локальных задач, когда LocalSet не запущен. Переданный future начнёт выполняться при следующем запуске LocalSet, даже если вы не ожидаете возвращённый 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;pub fn block_on<F>(&self, rt: &Runtime, future: F) -> F::Outputwhere F: Future,
Выполняет future до завершения в предоставленной среде выполнения, обслуживая в текущем потоке все локальные future, запущенные в этом наборе задач.
Этот метод выполняет переданный future в среде выполнения, блокируя выполнение до его завершения и возвращая полученный результат. Все задачи и таймеры, созданные future, будут выполняться в среде выполнения. Future также может вызвать spawn_local, чтобы spawn_local дополнительные локальные future в текущем потоке.
Этот метод нельзя вызывать из асинхронного контекста.
Панические завершения
Эта функция вызывает панику, если исполнитель достиг предельной нагрузки, если предоставленный future вызывает панику или если функция вызвана в контексте асинхронного выполнения.
Примечания
Поскольку эта функция внутри вызывает Runtime::block_on и обслуживает future из локального набора задач внутри этого вызова до block_on, локальные future не могут использовать блокировку на месте. Если из локальной задачи необходимо выполнить блокирующий вызов, вместо этого можно использовать API spawn_blocking.
Например, это вызовет панику:
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();
})
Однако это не вызовет панику:
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();
})pub async fn run_until<F>(&self, future: F) -> F::Outputwhere F: Future,
Выполняет future до завершения в локальном наборе и возвращает его результат.
Этот метод возвращает future, который выполняет переданный future с локальным набором, позволяя ему вызывать spawn_local для запуска дополнительных !Send future. Все локальные future, запущенные в локальном наборе, будут выполняться в фоновом режиме, пока не завершится future, переданный в run_until. Когда future, переданный в run_until, завершится, все незавершённые локальные future останутся в локальном наборе и будут выполняться при последующих вызовах run_until или при ожидании самого локального набора.
Безопасность при отмене
Этот метод безопасен при отмене, если future безопасен при отмене.
Примеры
use tokio::task;
task::LocalSet::new().run_until(async {
task::spawn_local(async move {
// ...
}).await.unwrap();
// ...
}).await;impl LocalSet
pub fn unhandled_panic(&mut self, behavior: UnhandledPanic) -> &mut Self
tokio_unstable.Настраивает поведение LocalSet при необработанной панике в запущенной задаче.
По умолчанию необработанная паника (то есть паника, не перехваченная с помощью std::panic::catch_unwind) не влияет на выполнение LocalSet. Значение ошибки паники передаётся в JoinHandle задачи, а все остальные запущенные задачи продолжают выполняться.
Параметр unhandled_panic позволяет настроить это поведение.
-
UnhandledPanic::Ignore— поведение по умолчанию. Паники в запущенных задачах не влияют на выполнениеLocalSet. -
UnhandledPanic::ShutdownRuntimeзаставитLocalSetнемедленно завершить работу при панике в запущенной задаче, даже если еёJoinHandleещё не был удалён. Все остальные запущенные задачи немедленно завершатся, а последующие вызовыLocalSet::block_onиLocalSet::run_untilвызовут панику.
Паники
Этот метод вызывает панику, если его вызвать после начала выполнения LocalSet.
Нестабильный API
Этот параметр в настоящее время является нестабильным, а его реализация не завершена. В будущем API может измениться или быть удалён. Подробнее см. tokio-rs/tokio#4516.
Примеры
В следующем примере показан LocalSet, настроенный на завершение работы при панике. Первая запущенная задача вызывает панику, в результате чего LocalSet завершает работу. Вторая запущенная задача не успевает выполниться. Вызов run_until вызовет панику, поскольку среда выполнения была принудительно остановлена.
use tokio::runtime::UnhandledPanic;
tokio::task::LocalSet::new()
.unhandled_panic(UnhandledPanic::ShutdownRuntime)
.run_until(async {
tokio::task::spawn_local(async { panic!("boom"); });
tokio::task::spawn_local(async {
// This task never completes
});
// Do some work, but `run_until` will panic before it completes
})
.await;
Реализации трейтов
Автоматические реализации трейтов
impl !Freeze for LocalSet
impl !RefUnwindSafe for LocalSet
impl !Send for LocalSet
impl !Sync for LocalSet
impl !UnwindSafe for LocalSet
impl Unpin for LocalSet
impl UnsafeUnpin for LocalSet
Общие реализации
impl<T> BorrowMut<T> for Twhere T: ?Sized,
fn borrow_mut(&mut self) -> &mut T
impl<T> Instrument for T
fn instrument(self, span: Span) -> Instrumented<Self> ⓘ
fn in_current_span(self) -> Instrumented<Self> ⓘ
impl<F> IntoFuture for Fwhere F: Future,
type Output = <F as Future>::Output
type IntoFuture = F
fn into_future(self) -> <F as IntoFuture>::IntoFuture
impl<T> WithSubscriber for T
fn with_subscriber<S>(self, subscriber: S) -> WithDispatch<Self> ⓘ
fn with_current_subscriber(self) -> WithDispatch<Self> ⓘ
Subscriber к этому типу, возвращая обёртку WithDispatch. Подробнее
MIT License
Copyright © Tokio Contributors
https://docs.rs/tokio/1.53.1/tokio/task/struct.LocalSet.html