Раздел 23 · Rust

Tokio изнутри

senior~50 мин

открытый урокЭтот раздел читается без входа. Войди, чтобы отмечать прогресс, вести заметки и решать задачи в редакторе. войти

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!.

Домашка