Tokio изнутри
открытый урокЭтот раздел читается без входа. Войди, чтобы отмечать прогресс, вести заметки и решать задачи в редакторе. войти
Tokio изнутри
Прошлый урок собрал однопоточный executor: одна очередь, один поток крутит
poll. Этого хватает для учёбы, но не масштабируется: восемь ядер простаивают, пока один поток тащит все задачи. Tokio решает это многопоточным планировщиком с work-stealing. Сегодня разбираем его устройство: локальные и глобальные очереди, воровство задач, I/O-драйвер, таймерное колесо и кооперативную вытесняемость. Свой однопоточный рантайм у тебя уже есть, теперь смотрим, что добавляет промышленный.
Зачем не одна очередь на всех
Очевидное решение для многих потоков это одна общая очередь задач под Mutex: все воркеры берут из неё. И оно плохо масштабируется. Каждое взятие задачи это захват общего замка, и на восьми ядрах потоки больше дерутся за Mutex, чем работают. Это называется contention, и из-за него наивная схема упирается в потолок.
Tokio устроен иначе. У каждого воркера своя локальная очередь, к которой он обращается почти без синхронизации. Есть одна глобальная очередь-инжектор для задач, заспавненных снаружи. А когда у воркера кончились свои задачи, он не ждёт, а идёт воровать чужие.
Work-stealing
Work-stealing отвечает на вопрос: где простаивающий воркер берёт работу. Алгоритм каждого воркера такой: сначала своя локальная очередь (самый дешёвый путь, данные в кэше), если пусто, то пачка из глобального инжектора, если и там пусто, то воровство у соседей. Соберём эту геометрию поверх crossbeam-deque: Worker это локальная очередь, Stealer это ручка для воровства из неё, Injector это глобальная очередь.
use crossbeam_deque::{Injector, Stealer, Worker, Steal};
/// Найти задачу: сначала своя локальная, потом инжектор, потом соседи.
fn find_job(index: usize, local: &Worker<Job>, pool: &Pool) -> Option<Job> {
// 1. Локальная очередь: самый дешёвый путь, без воровства.
if let Some(job) = local.pop() {
return Some(job);
}
loop {
// 2. Глобальный инжектор: тащим пачку в локальную очередь разом.
match pool.injector.steal_batch_and_pop(local) {
Steal::Success(job) => return Some(job),
Steal::Retry => continue,
Steal::Empty => {}
}
// 3. Воровство у соседей: обходим чужие Stealer'ы, кроме своего.
let mut retry = false;
for (i, stealer) in pool.stealers.iter().enumerate() {
if i == index {
continue;
}
match stealer.steal() {
Steal::Success(job) => return Some(job),
Steal::Retry => retry = true,
Steal::Empty => {}
}
}
if !retry {
return None; // работы нет нигде
}
}
}
Обрати внимание на steal_batch_and_pop: из инжектора воркер тащит не одну задачу, а сразу пачку себе в локальную очередь. Это снижает походы к общему ресурсу: украл пачку один раз, дальше работаешь со своей очередью без синхронизации. Реальный Tokio так же ворует у соседей половину их очереди, а не по одной задаче, по той же причине.
Цикл воркера крутит find_job и выполняет найденное; если работы нет нигде, паркуется до новой задачи:
fn worker_loop(index: usize, local: Worker<Job>, pool: Arc<Pool>) {
loop {
if let Some(job) = find_job(index, &local, &pool) {
job();
continue;
}
if pool.shutdown.load(Ordering::SeqCst) {
return;
}
// Работы нет: паркуемся до новой задачи (с коротким таймаутом, чтобы
// не проспать задачу, попавшую между find_job и парковкой).
let guard = pool.idle.lock().unwrap();
let _ = pool.idle_cv.wait_timeout(guard, Duration::from_millis(1));
}
}
Наш учебный пул гоняет синхронные замыкания FnOnce, а не футуры: урок про геометрию очередей и распределение работы, а не про poll. Настоящий Tokio в этих же очередях держит задачи-футуры, но расстановка локальная-глобальная-воровство ровно эта. Кстати, ту же задачу M:N-планирования (много зелёных задач на немногих системных потоках) решает планировщик горутин в Go, разобранный в уроке про конкурентность вглубь; идея локальных очередей и воровства там та же.
Сборка пула
Осталось показать, как очереди заводятся и как пул запускается. Геометрия вот в чём: локальные очереди (Worker) создаются заранее, чтобы с каждой снять ручку воровства (Stealer) и сложить все ручки в общий Pool. Тогда любой воркер видит Stealer’ы всех соседей.
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
use std::sync::{Arc, Condvar, Mutex};
use std::thread::{self, JoinHandle};
use crossbeam_deque::{Injector, Stealer, Worker};
type Job = Box<dyn FnOnce() + Send + 'static>;
struct Pool {
injector: Injector<Job>, // глобальная очередь: сюда кладёт spawn
stealers: Vec<Stealer<Job>>, // ручки воровства у всех локальных очередей
pending: AtomicUsize, // сколько задач ещё не завершилось
shutdown: AtomicBool, // сигнал остановки воркерам
idle: Mutex<()>, // воркеры паркуются здесь, когда работы нет
idle_cv: Condvar,
done: Mutex<bool>, // кого будить, когда pending дошёл до нуля
done_cv: Condvar,
}
pub struct WorkStealingPool {
pool: Arc<Pool>,
workers: Vec<JoinHandle<()>>,
}
impl WorkStealingPool {
pub fn new(workers: usize) -> Self {
assert!(workers > 0, "нужен хотя бы один воркер");
// Сначала локальные очереди, потом снимаем с каждой Stealer.
let local_queues: Vec<Worker<Job>> =
(0..workers).map(|_| Worker::new_fifo()).collect();
let stealers = local_queues.iter().map(|w| w.stealer()).collect();
let pool = Arc::new(Pool {
injector: Injector::new(),
stealers,
pending: AtomicUsize::new(0),
shutdown: AtomicBool::new(false),
idle: Mutex::new(()),
idle_cv: Condvar::new(),
done: Mutex::new(false),
done_cv: Condvar::new(),
});
// Каждой локальной очереди свой поток с циклом worker_loop.
let handles = local_queues
.into_iter()
.enumerate()
.map(|(index, local)| {
let pool = pool.clone();
thread::spawn(move || worker_loop(index, local, pool))
})
.collect();
WorkStealingPool { pool, workers: handles }
}
pub fn spawn(&self, job: impl FnOnce() + Send + 'static) {
self.pool.pending.fetch_add(1, Ordering::SeqCst);
self.pool.injector.push(Box::new(job)); // внешний spawn -> в инжектор
self.pool.idle_cv.notify_all(); // разбудить спящих воркеров
}
pub fn run_until_complete(&self) {
let mut done = self.pool.done.lock().unwrap();
while self.pool.pending.load(Ordering::SeqCst) > 0 {
done = self.pool.done_cv.wait(done).unwrap();
}
*done = false;
}
}
impl Drop for WorkStealingPool {
fn drop(&mut self) {
self.pool.shutdown.store(true, Ordering::SeqCst); // сигналим остановку
self.pool.idle_cv.notify_all(); // будим спящих
for handle in self.workers.drain(..) {
let _ = handle.join(); // ждём завершения
}
}
}
Чтобы run_until_complete знал, когда всё готово, воркер после каждой задачи уменьшает pending, и когда счётчик дошёл до нуля, будит ожидающего. К циклу воркера из прошлого раздела добавляется ровно этот хвост:
if let Some(job) = find_job(index, &local, &pool) {
job();
// Задача готова: декремент pending, и если это была последняя, будим
// того, кто ждёт в run_until_complete.
if pool.pending.fetch_sub(1, Ordering::SeqCst) == 1 {
let mut done = pool.done.lock().unwrap();
*done = true;
pool.done_cv.notify_all();
}
continue;
}
Запуск в три шага: создал пул, накидал задач, дождался.
let pool = WorkStealingPool::new(4);
let counter = Arc::new(AtomicUsize::new(0));
for _ in 0..1000 {
let counter = counter.clone();
pool.spawn(move || { counter.fetch_add(1, Ordering::SeqCst); });
}
pool.run_until_complete();
assert_eq!(counter.load(Ordering::SeqCst), 1000);
Тысяча задач разойдётся по четырём потокам сама: кто освободился, тот и своровал. Никакого центрального диспетчера, только локальные очереди и find_job.
Чего ещё нет в нашей схеме
Учебный пул показывает скелет. Промышленный планировщик добавляет к нему несколько важных деталей.
LIFO-слот. Когда задача будит другую (например, послала в канал и разбудила приёмника), Tokio кладёт разбуженную не в хвост очереди, а в особый LIFO-слот, чтобы выполнить её следующей. Логика: эти две задачи общаются, их данные горячие в кэше, выгодно дать им поработать подряд. Это оптимизация под типичный паттерн «запрос-ответ».
I/O-драйвер. Это reactor из прошлого урока, поверх mio. Один из потоков (или выделенный) крутит epoll и будит задачи, ждущие сокетов. Связь с executor’ом ровно та, что мы собрали руками: Token -> Waker.
Таймерное колесо. Чтобы поддержать tokio::time::sleep и таймауты, рантайму нужно эффективно хранить тысячи таймеров и быстро находить сработавшие. Линейный список не годится. Используется таймерное колесо: кольцевой массив корзин, где таймер кладётся в корзину по своему сроку, а продвижение времени это переход к следующей корзине. Вставка и снятие за O(1). Наш самодельный таймер на линейном списке делал бы O(n) на каждый тик и умер бы на больших объёмах.
Где ещё живут те же кубики
Чтобы карта была полной, оглянись на соседей по экосистеме. smol устроен идейно очень близко к тому, что мы собрали руками: executor с очередью задач плюс reactor, его исходники читаются за вечер. glommio это другая архитектура, thread-per-core поверх io_uring: у каждого ядра свой reactor, задачи между потоками не переезжают. А сам Tokio на свежих ядрах Linux умеет брать io_uring как бэкенд I/O-драйвера вместо epoll. Кубики везде те же три; меняется лишь то, как их раскладывают по потокам и ядрам.
Кооперативная вытесняемость
И последняя деталь, без которой большой рантайм заклинивает. Async в Rust кооперативный: задача отдаёт управление только на await. Если внутри задачи крутится тяжёлый цикл без единого await (или есть await, который всегда сразу Ready), задача никогда не возвращает управление и монополизирует воркер. Соседние задачи на этом потоке стоят.
Tokio защищается бюджетом: у каждой задачи есть счётчик, и после примерно 128 готовых операций над ресурсами Tokio (чтений из канала, из сокета) задача принудительно возвращает Pending, отдавая ход соседям, даже если данные ещё есть. Это компромисс: чистая кооперативность хрупка (одна жадная задача вешает воркер), а полноценное вытеснение по таймеру, как у потоков ОС, для async-задач не сделать. Разница вытесняющей и кооперативной моделей подробно разбиралась в уроке про гранулярность.
Два режима рантайма
Tokio собирается в двух режимах, и выбор между ними это выбор «нужен ли мне тот work-stealing, что мы только что собрали».
Многопоточный рантайм (по умолчанию) это ровно наш пул: несколько воркеров, локальные очереди, воровство. Задачи в нём должны быть Send, ведь воровство перекидывает их между потоками.
Однопоточный рантайм (current_thread) это наш block_on-executor из урока 51: одна очередь, один поток, никакого воровства. Задачам не нужен Send, накладных расходов на синхронизацию очередей нет. Это выбор для CLI, тестов и случаев, где параллелизм по ядрам не нужен.
// Многопоточный (по умолчанию). worker_threads задаёт число воркеров.
#[tokio::main(flavor = "multi_thread", worker_threads = 8)]
async fn main() { /* ... */ }
// Однопоточный: один поток крутит всё, задачи могут быть !Send.
#[tokio::main(flavor = "current_thread")]
async fn main() { /* ... */ }
// То же руками через Builder, когда нужен контроль над рантаймом:
let rt = tokio::runtime::Builder::new_multi_thread()
.worker_threads(8)
.enable_all() // включить I/O-драйвер и таймеры
.build()
.unwrap();
rt.block_on(async { /* ... */ });
spawn_blocking: тяжёлый синхронный код
Кооперативный бюджет спасает от жадной задачи, которая часто трогает ресурсы Tokio. Но он бессилен против честно блокирующего кода: хеширование мегабайта, синхронный std::fs::read, парсинг без единого await. Такой код вообще не доходит до точки отдачи и держит воркер мёртвой хваткой, пока не закончит. На многопоточном рантайме это съедает одно ядро из восьми, на однопоточном вешает вообще всё.
Решение это spawn_blocking: отправить блокирующую работу на отдельный пул потоков, который Tokio держит именно под это (по умолчанию до 512 потоков). Async-воркеры остаются свободны, а вызывающая задача .await-ит результат через JoinHandle, как любую другую.
// Тяжёлый синхронный хеш уходит на блокирующий пул, async-воркеры свободны.
let digest = tokio::task::spawn_blocking(move || expensive_hash(&data))
.await
.expect("blocking task panicked");
Правило практическое: если кусок кода может работать дольше пары десятков микросекунд и не содержит await, ему место в spawn_blocking, а не в async-задаче. Это та же мысль, что «не блокируй event loop» из мира JS, только тут у тебя есть честный пул, куда блокирующее можно отселить.
Что унести из урока
Многопоточный планировщик Tokio уходит от одной общей очереди под Mutex, потому что она тонет в contention. Вместо этого у каждого воркера своя локальная очередь без синхронизации, есть глобальный инжектор для внешних задач, а простаивающий воркер ворует работу: сначала пачку из инжектора, потом у соседей. Поверх этого скелета лежат LIFO-слот для общающихся задач, I/O-драйвер на mio, таймерное колесо для O(1) таймеров и кооперативный бюджет, который не даёт жадной задаче подвесить воркер. Рантайм собирается в двух режимах: многопоточный (наш work-stealing пул, задачи Send) и current_thread (наш block_on-executor, задачи могут быть !Send). Async кооперативный, поэтому за справедливостью следит бюджет, а честно блокирующий код выносят на отдельный пул через spawn_blocking.
Дальше уходим от планирования к семантике: что значит отменить async-задачу, почему отмена это просто drop, и какие ловушки в select!.