Rust Асинхронный Rust: future, tokio, async/await и его подводные камни
0%

Асинхронный Rust: future, tokio, async/await и его подводные камни

Асинхронный Rust: future, tokio, async/await и его подводные камни

В предыдущей статье мы разобрали конкурентность на потоках: std::thread::spawn, каналы, Arc<Mutex<T>>, типажи Send и Sync. Этой модели хватает для огромного класса задач, и начинать почти всегда стоит с неё.

Но есть задача, на которой она ломается. Вот она, в самой честной формулировке: сервер держит 50 000 открытых соединений, из которых в любой момент что-то передают человек двести. Остальные 49 800 просто висят — WebSocket-подписки, keep-alive, долгие опросы. Если каждому соединению выдать поток ОС, вы платите за 50 000 стеков и за то, что планировщик ядра перебирает 50 000 сущностей, 99% которых спят.

поток ОС на Linux            задача tokio
────────────────────────     ──────────────────────────
8 МиБ виртуального стека     размер автомата состояний: обычно 100 байт – 2 КиБ
~8–16 КиБ реального RSS      +64 байта служебных данных задачи
переключение ~1–2 мкс        переключение = возврат из poll, десятки наносекунд
создание ~10–20 мкс          создание = одна аллокация
предел практический ~10⁴     предел практический ~10⁶

Асинхронность — это способ перестать платить за ожидание. Не за вычисления: async не делает вашу программу быстрее ни на один такт, он лишь позволяет не занимать поток, пока данные не пришли. Это статья о том, как именно Rust это делает, почему устройство именно такое, и где новичок обязательно порежется.

Сразу зафиксируем главное отличие от всех знакомых языков, потому что из него вырастут все грабли: future в Rust инертна. async fn не запускает ничего. Он возвращает значение, которое ничего не делает, пока кто-то извне не вызовет у него poll. Ни Promise из JavaScript, ни Task из C# (трек csharp), ни горутина (трек golang) так себя не ведут: они стартуют сами.

Future: трейт из трёх строк

Вся асинхронность Rust стоит на одном типаже из стандартной библиотеки (std::future::Future):

pub trait Future {
    type Output;
    fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output>;
}

pub enum Poll<T> {
    Ready(T),   // результат готов, больше меня не опрашивай
    Pending,    // я не готов; когда буду — разбужу через Waker из cx
}

Прочитайте poll как вопрос: «ты уже закончил?» Ответ либо «да, вот значение», либо «нет, я сохранил твой будильник и позвоню». Никакого блокирования внутри poll быть не должно — это фундаментальный контракт, и его нарушение и есть грабля номер один.

Напишем future руками, чтобы контракт стал осязаемым:

use std::future::Future;
use std::pin::Pin;
use std::task::{Context, Poll};
use std::time::Instant;

/// Готова, когда наступил дедлайн.
struct Delay {
    deadline: Instant,
}

impl Future for Delay {
    type Output = ();

    fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<()> {
        if Instant::now() >= self.deadline {
            Poll::Ready(())
        } else {
            // Наивная реализация: просим опросить нас снова прямо сейчас.
            // Это busy-loop и так делать нельзя, зато видно суть контракта:
            // "Pending" обязан сопровождаться обещанием разбудить.
            cx.waker().wake_by_ref();
            Poll::Pending
        }
    }
}

Настоящий tokio::time::sleep отличается ровно одним: он не крутит процессор, а регистрирует Waker в колесе таймеров рантайма и возвращает Pending один раз. Драйвер времени вызовет wake(), когда срок придёт. Та же логика у сокета: Waker уходит в epoll, и задачу разбудят по готовности файлового дескриптора (механику epoll подробно разбирает «Ввод-вывод и системные вызовы»).

async fn — это автомат состояний

Компилятор превращает каждый async fn и каждый async блок в анонимный тип, реализующий Future. Никакой магии там нет: это тот же приём, что в C# со Task и в генераторах Python, только с точностью до байта.

async fn handle(id: u64) -> usize {
    let row = load(id).await;
    let name = &row.name;
    let n = store(name).await;
    n + row.name.len()
}

Компилятор строит из этого нечто по смыслу эквивалентное вот такому перечислению (реальный код — корутина, но форма именно эта):

enum HandleFuture {
    Start { id: u64 },
    AwaitingLoad { fut: Load },
    AwaitingStore { fut: Store, row: Row },   // row живёт через await, значит хранится тут
    Done,
}

Из этой картинки прямо читаются три вещи, которые придётся принять:

  1. Локальная переменная, живущая через .await, переезжает в объект future. Не живущая — остаётся в обычном кадре стека на время одного poll. Отсюда правило «не держите большие буферы через await, если их можно не держать».
  2. Размер future — примерно максимум по состояниям, а не сумма. Но вложенные future складываются вглубь, поэтому глубокая цепочка async fn даёт неприятно большой объект.
  3. &row.name — указатель внутрь самого объекта future. То есть future самоссылающаяся, а значит её нельзя двигать после начала работы. Именно поэтому poll принимает Pin<&mut Self>, а не &mut Self.

Автомат состояний future в памяти и зачем нужен Pin

Pin<P> — это обещание системе типов: «пока значение живо, оно не переедет». Из Pin<&mut T> нельзя получить &mut T, если T: !Unpin, а все future, сгенерированные из async, именно !Unpin. Двигать их можно свободно до первого poll (поэтому let f = foo(); bar(f); работает), а после — адрес зафиксирован. Практический вывод: если компилятор просит Pin, используйте Box::pin (аллокация, зато 'static), std::pin::pin! или tokio::pin! (на стеке, без аллокации). Детали — в документации std::pin, а самое понятное объяснение с картинками — в Async Book, глава про Pinning.

Подробнее о том, почему &mut T инвариантен и как связаны Pin и модель заимствования, — во «Временах жизни».

Кто вызывает poll: исполнитель, реактор, Waker

Стандартная библиотека даёт трейт Future, Waker и Context — и всё. Ни исполнителя, ни таймеров, ни сокетов в std нет. Это осознанное решение: асинхронный рантайм должен быть заменяемым, потому что требования у веб-сервера, у прошивки микроконтроллера и у ядра ОС разные.

Внутри рантайма всегда две половины:

  • Исполнитель (executor) — крутит цикл «взять задачу из очереди → вызвать poll → если Pending, забыть о ней».
  • Реактор (driver) — общается с ядром: один epoll_wait на тысячи сокетов, колесо таймеров. Он же владеет Waker-ами и дёргает их по готовности.

Исполнитель писать несложно, и один раз это стоит сделать руками. Вот полноценный block_on на 35 строк — по сути это и есть futures::executor::block_on:

use std::future::Future;
use std::pin::pin;
use std::sync::{Arc, Condvar, Mutex};
use std::task::{Context, Poll, Wake, Waker};

/// Сигнал «меня разбудили»: флаг под мьютексом плюс условная переменная.
struct Signal {
    woken: Mutex<bool>,
    cv: Condvar,
}

impl Signal {
    fn wait(&self) {
        let mut woken = self.woken.lock().unwrap();
        while !*woken {
            woken = self.cv.wait(woken).unwrap();   // спим, не тратя процессор
        }
        *woken = false;                            // готовимся к следующему кругу
    }
    fn notify(&self) {
        let mut woken = self.woken.lock().unwrap();
        *woken = true;
        self.cv.notify_one();
    }
}

// Весь контракт Waker: «сделай так, чтобы меня опросили снова».
impl Wake for Signal {
    fn wake(self: Arc<Self>) { self.notify() }
    fn wake_by_ref(self: &Arc<Self>) { self.notify() }
}

fn block_on<F: Future>(future: F) -> F::Output {
    let mut future = pin!(future);          // прибиваем адрес: poll требует Pin
    let signal = Arc::new(Signal { woken: Mutex::new(false), cv: Condvar::new() });
    let waker = Waker::from(Arc::clone(&signal));
    let mut cx = Context::from_waker(&waker);

    loop {
        match future.as_mut().poll(&mut cx) {
            Poll::Ready(value) => return value,
            Poll::Pending => signal.wait(),
        }
    }
}

Тридцать пять строк — и у вас однозадачный рантайм. Разница между ним и tokio — не в идее, а в том, что tokio умеет много задач, воровать работу между потоками и разговаривать с epoll. Если хочется полного разбора с нуля — «Futures Explained in 200 Lines of Rust» и глава «Async/Await» в «Writing an OS in Rust», где исполнитель пишется прямо в ядре без стандартной библиотеки.

Жизненный цикл задачи

Задача (task) — это future, отданная рантайму во владение через spawn. Её состояния стоит держать в голове целиком, потому что половина ошибок — это неучтённый переход в «Отменена».

Из этой диаграммы сразу три неочевидных факта. Первый: Создана --> Отменена — если future не дождались, она не выполнится совсем, тихо и без ошибки. Второй: отмена возможна только в точках приостановки — код между двумя .await неотменяем, ждать придётся до конца. Третий: отмена — это Drop, поэтому у неё нет асинхронной фазы, и «закрыть соединение красиво при отмене» в общем случае невозможно (асинхронного Drop в языке нет — это признанное ограничение, см. wg-async).

tokio: рантайм, который вы будете использовать

На практике выбор невелик: tokio занимает подавляющую долю экосистемы, потому что на нём стоят hyper, axum, tonic, reqwest, sqlx. Альтернативы имеют смысл в узких нишах: smol — простота и малый размер, glommio и monoio — thread-per-core на io_uring, embassy — микроконтроллеры без ОС. Обзор состояния — «The State of Async Rust: Runtimes».

#[tokio::main]                            // сахар, разворачивающийся в код ниже
async fn main() {
    println!("привет из рантайма");
}

// Ровно то же самое руками — полезно видеть, что рантайм это просто объект:
fn main() {
    tokio::runtime::Builder::new_multi_thread()
        .worker_threads(4)                // по умолчанию = число логических ядер
        .enable_all()                     // включить драйверы ввода-вывода и времени
        .build()
        .unwrap()
        .block_on(async {
            println!("привет из рантайма");
        });
}

Топология рантайма tokio: очереди, воровство работы, драйверы, блокирующий пул

Самая важная деталь этой схемы: воркеров ровно столько, сколько ядер, и они не вытесняются. В Go рантайм с версии 1.14 умеет асинхронное вытеснение горутин, в Elixir на BEAM процессы вытесняются по счётчику редукций (трек elixir). В Rust точка переключения одна — возврат из poll, то есть .await. Задача, которая молотит процессор без await, держит воркер целиком и отъедает 1/N пропускной способности процесса. Это не баг, а прямое следствие того, что задача — stackless-корутина без своего стека и без возможности прервать её посередине. Как ядро планирует сами рабочие потоки — в «Процессы и планирование».

Второй важный нюанс — spawn требует Send + 'static, потому что work stealing может увезти задачу на другой поток в любой точке приостановки. Для !Send future есть tokio::task::LocalSet и spawn_local либо flavor = "current_thread".

Практика: сервер, который держит десятки тысяч соединений

use tokio::io::{AsyncReadExt, AsyncWriteExt};
use tokio::net::TcpListener;

#[tokio::main]
async fn main() -> std::io::Result<()> {
    let listener = TcpListener::bind("127.0.0.1:8080").await?;
    println!("слушаю 127.0.0.1:8080");

    loop {
        let (mut socket, peer) = listener.accept().await?;

        // Соединение = задача, а не поток: десятки байт вместо мегабайт стека.
        tokio::spawn(async move {
            let mut buf = [0u8; 4096];
            loop {
                match socket.read(&mut buf).await {
                    Ok(0) => break,                              // клиент закрыл
                    Ok(n) => {
                        if socket.write_all(&buf[..n]).await.is_err() {
                            break;                               // соединение оборвалось
                        }
                    }
                    Err(e) => { eprintln!("{peer}: ошибка чтения: {e}"); break; }
                }
            }
        });
    }
}

Обратите внимание: буфер [0u8; 4096] объявлен внутри задачи и живёт через await, значит он часть автомата состояний. 50 000 соединений — это 200 МиБ только на буферы. На больших числах соединений буферы берут из пула (bytes::BytesMut) или уменьшают.

Теперь про разницу между конкурентностью и параллелизмом, которую в async путают чаще всего:

use tokio::time::{sleep, Duration, Instant};

async fn fetch(name: &str, ms: u64) -> String {
    sleep(Duration::from_millis(ms)).await;
    format!("{name} готов")
}

#[tokio::main]
async fn main() {
    let t0 = Instant::now();
    let _a = fetch("a", 100).await;          // последовательно:
    let _b = fetch("b", 150).await;          // каждая ждёт предыдущую
    let _c = fetch("c", 200).await;
    println!("последовательно: {:?}", t0.elapsed());

    let t1 = Instant::now();
    // join! опрашивает три future в ОДНОЙ задаче, на ОДНОМ потоке.
    // Это конкурентность без параллелизма — и её достаточно для ожидания.
    let (_a, _b, _c) = tokio::join!(fetch("a", 100), fetch("b", 150), fetch("c", 200));
    println!("join!:           {:?}", t1.elapsed());
}
последовательно: 453.812ms
join!:           201.447ms

join! — конкурентность внутри одной задачи. spawn — конкурентность плюс возможный параллелизм на разных воркерах, ценой требования Send + 'static и одной аллокации. Правило: если работа только ждёт — join!; если ещё и считает — spawn.

Для набора однотипных задач с ожиданием результатов используйте JoinSet: он умеет отдавать результаты по мере готовности и, что важнее, при уронении отменяет все свои задачи — это ближайшее к структурной конкурентности, что есть в tokio из коробки.

use tokio::task::JoinSet;

#[tokio::main]
async fn main() {
    let mut set = JoinSet::new();
    for id in 0..100u64 {
        set.spawn(async move { (id, work(id).await) });
    }

    while let Some(res) = set.join_next().await {
        match res {
            Ok((id, value)) => println!("{id} -> {value}"),
            // Паника внутри задачи не роняет процесс: она приезжает как JoinError.
            Err(e) if e.is_panic() => eprintln!("задача паниковала: {e}"),
            Err(e) => eprintln!("задача отменена: {e}"),
        }
    }
}

Восемь грабель, на которые наступают все

1. Блокирующий вызов внутри async

Это ошибка номер один с большим отрывом. poll обязан вернуться быстро; всё, что занимает поток надолго, останавливает не одну задачу, а весь воркер.

#[tokio::main(flavor = "current_thread")]
async fn main() {
    let mut handles = Vec::new();
    for i in 0..4 {
        handles.push(tokio::spawn(async move {
            // ГРАБЛЯ: sleep из std блокирует ПОТОК, а не задачу.
            std::thread::sleep(std::time::Duration::from_secs(1));
            println!("задача {i} готова");
        }));
    }
    for h in handles { h.await.unwrap(); }
}
задача 0 готова     # через 1 с
задача 1 готова     # через 2 с
задача 2 готова     # через 3 с
задача 3 готова     # через 4 с — всего 4 секунды вместо одной

Замените на tokio::time::sleep().await — получите одну секунду и произвольный порядок. Коварство в том, что на многопоточном рантайме с 16 воркерами такая ошибка не видна на малой нагрузке: 16 блокирующих вызовов проходят «параллельно», деградация появляется только в проде под трафиком, и выглядит как необъяснимо огромный p99.

Список запрещённого в асинхронном коде: std::fs, std::net, std::thread::sleep, Mutex::lock из std под настоящей конкуренцией, синхронные драйверы СУБД (diesel, postgres, rusqlite), FFI-вызовы неизвестной длительности, println! в цикле (запись в stdout блокирующая), тяжёлый CPU — хеширование паролей, сжатие, обработка картинок, разбор мегабайтных JSON.

// Правильно: тяжёлая синхронная работа уезжает в отдельный пул.
let hash = tokio::task::spawn_blocking(move || {
    argon2_hash(&password)     // 200 мс чистого CPU — воркеру такого нельзя
}).await?;

Определение слова «блокирует» с примерами и разбором границ — Alice Ryhl, «Async: What is blocking?», обязательное чтение. Как измерять последствия — в «Конкурентность и производительность».

2. Send, и почему std::sync::MutexGuard нельзя тащить через await

use std::sync::{Arc, Mutex};

async fn fetch_delta() -> u64 { 1 }

fn broken(counter: Arc<Mutex<u64>>) {
    tokio::spawn(async move {
        let mut guard = counter.lock().unwrap();
        *guard += fetch_delta().await;     // guard жив через await
    });
}
error: future cannot be sent between threads safely
   --> src/main.rs:6:18
    |
6   |     tokio::spawn(async move {
    |                  ^^^^^^^^^^ future created by async block is not `Send`
    |
note: future is not `Send` as this value is used across an await
   --> src/main.rs:8:29
    |
7   |         let mut guard = counter.lock().unwrap();
    |             --------- has type `std::sync::MutexGuard<'_, u64>` which is not `Send`
8   |         *guard += fetch_delta().await;
    |                                 ^^^^^ await occurs here, with `guard` maybe used later
note: required by a bound in `tokio::spawn`

Читайте это сообщение как рассказ в три акта: что не Send (сам блок), почему (конкретное значение живёт через await), и кто это потребовал (tokio::spawn). Причина реальна: MutexGuard в POSIX привязан к потоку, а work stealing может продолжить задачу на другом воркере — освобождение мьютекса из чужого потока это undefined behavior (детали UB).

Два варианта починки, и первый почти всегда правильный:

// Вариант A: сузить область блокировки так, чтобы await был снаружи. Дёшево и быстро.
let delta = fetch_delta().await;
*counter.lock().unwrap() += delta;

// Вариант B: асинхронный мьютекс — только если блокировку ДЕЙСТВИТЕЛЬНО надо
// держать через await (например, вы владеете соединением и пишете в него по частям).
use tokio::sync::Mutex;
let mut guard = counter.lock().await;   // guard: Send, задача может мигрировать
*guard += fetch_delta().await;

Официальная рекомендация tokio: по умолчанию берите std::sync::Mutex, он в разы дешевле; tokio::sync::Mutex — только когда критическая секция содержит await. И помните: асинхронный мьютекс не спасает от логической проблемы «все задачи стоят в очереди за одной блокировкой».

3. Отмена: то, чего нет ни в Go, ни в C#

select! опрашивает несколько future и, как только одна готова, роняет остальные. Отсюда понятие cancellation safety: безопасна ли future к тому, что её уронят посередине.

use tokio::io::{AsyncBufReadExt, BufReader};

// НЕПРАВИЛЬНО: read_line НЕ cancel-safe.
loop {
    let mut line = String::new();
    tokio::select! {
        res = reader.read_line(&mut line) => handle(res, &line).await,
        _ = token.cancelled() => break,
    }
}

Если сработает вторая ветка, read_line уронят — а байты из сокета уже прочитаны в её внутренний буфер, и он исчезнет вместе с future. Данные потеряны молча, без ошибки. Такой баг воспроизводится раз в сутки на проде и не воспроизводится в тестах никогда.

Три рабочих подхода:

// A. Создать future ОДИН раз и прикрепить: между итерациями она не пересоздаётся.
let read = reader.read_line(&mut line);
tokio::pin!(read);
loop {
    tokio::select! {
        res = &mut read => { /* ... */ }
        _ = token.cancelled() => break,
    }
}

// B. Изолировать ввод-вывод в своей задаче и общаться каналом:
//    mpsc::Receiver::recv cancel-safe по документации.
tokio::select! {
    Some(line) = rx.recv() => handle(line).await,
    _ = token.cancelled() => return,
}

Третий подход — просто читать документацию: в tokio у каждого метода есть раздел «Cancel safety». Cancel-safe: mpsc::Receiver::recv, broadcast::Receiver::recv, Mutex::lock, AsyncReadExt::read, TcpListener::accept. Не cancel-safe: AsyncBufReadExt::read_line, AsyncWriteExt::write_all, read_exact. Точный перечень — в документации tokio::select!.

Отдельная ловушка отмены: уронить JoinHandle НЕ отменяет задачу. Задача становится отсоединённой и живёт своей жизнью. Отменяет только handle.abort() или уронение всего JoinSet. Обратная ловушка: tokio::time::timeout(d, fut) отменяет fut по истечении срока, и если внутри была транзакция в БД — она обрывается в произвольной точке приостановки.

Грамотное завершение делают через CancellationToken из tokio-util плюс JoinSet, паттерн описан в tokio topics: Graceful Shutdown.

4. Future ничего не делает, пока её не опросят

async fn save(record: &Record) -> Result<(), Error> { /* ... */ }

async fn handler(r: Record) {
    save(&r);          // забыли .await — не сохранилось НИЧЕГО
    respond_ok().await;
}
warning: unused implementer of `Future` that must be used
 --> src/main.rs:4:5
  |
4 |     save(&r);
  |     ^^^^^^^^
  |
  = note: futures do nothing unless you `.await` or poll them

Компилятор предупреждает, но это warning, а не ошибка, и в шуме сборки он теряется. Лечится #![deny(unused_must_use)] в корне крейта. Симметричный вариант той же ошибки — накопить Vec<impl Future> и не дождаться его: задачи не выполнятся. В C# Task уже запущен и «забыл await» означает лишь потерю результата и проглоченное исключение; в Rust не происходит вообще ничего.

5. async в трейтах: стабилизировали, но не до конца

С Rust 1.75 async fn в трейтах компилируется (анонс):

trait Storage {
    async fn get(&self, key: &str) -> Option<Vec<u8>>;
}

А вот попытка отдать это в spawn — нет:

fn spawn_it<S: Storage + Send + Sync + 'static>(s: S) {
    tokio::spawn(async move { s.get("k").await; });
}
error: future cannot be sent between threads safely
note: future is not `Send` as it awaits another future which is not `Send`

Проблема в том, что на стабильном Rust нельзя написать границу «возвращаемая трейтом future должна быть Send». Плюс такой трейт не dyn-совместим: Box<dyn Storage> не соберётся. Практические выходы: #[async_trait] (крейт async-trait, боксирует каждый вызов — одна аллокация, зато работает всё), trait_variant::make для генерации Send-варианта трейта, либо вручную объявить fn get(&self) -> Pin<Box<dyn Future<Output = ...> + Send + '_>>. Это до сих пор самый шероховатый угол асинхронного Rust, и о нём стоит знать до того, как вы спроектируете абстракцию над хранилищем.

6. Рантайм не запущен — или запущен дважды

thread 'main' panicked at 'there is no reactor running,
must be called from the context of a Tokio 1.x runtime'

Это tokio::spawn, TcpStream::connect или sleep вне block_on. Типично в тестах — забыли #[tokio::test], — и в коде, вызванном из C через FFI. Обратная ситуация:

thread 'tokio-runtime-worker' panicked at 'Cannot start a runtime from within a runtime.
This happens because a function (like `block_on`) attempted to block the current thread
while the thread is being used to drive asynchronous tasks.'

Так падает попытка вызвать block_on (или синхронную обёртку над ним из чужой библиотеки) внутри асинхронной функции. Правильный мост в обе стороны — Handle::current().block_on только на потоке вне рантайма, spawn_blocking для синхронного кода, Handle::block_on изнутри spawn_blocking; см. tokio topics: Bridging with sync code. Родственная беда — смешение рантаймов: reqwest требует tokio, и на async-std он тихо не работает или паникует. Ещё одна причина, по которой выбор рантайма делается один раз на весь проект.

7. Размер future и рекурсия

async fn walk(dir: PathBuf) -> u64 {
    let mut total = 0;
    let mut rd = tokio::fs::read_dir(&dir).await.unwrap();
    while let Some(e) = rd.next_entry().await.unwrap() {
        total += walk(e.path()).await;   // рекурсия
    }
    total
}
error[E0733]: recursion in an async fn requires boxing
 --> src/main.rs:1:1
  |
1 | async fn walk(dir: PathBuf) -> u64 {
  | ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
  |
note: a recursive `async fn` must introduce indirection such as `Box::pin` to avoid
      an infinitely sized future

Причина прозрачна из модели: размер автомата состояний считается статически, а рекурсия делает его бесконечным. Починка — разрыв через указатель:

fn walk(dir: PathBuf) -> Pin<Box<dyn Future<Output = u64> + Send>> {
    Box::pin(async move { /* тело как было */ })
}

Начиная с Rust 1.77 боксированная рекурсия работает без внешних крейтов; раньше требовался async_recursion. Родственная проблема без сообщения об ошибке — раздутая future: глубокая цепочка async fn, каждая со своими локальными переменными, даёт объект на десятки килобайт, который копируется при каждом spawn и способен переполнить стек. Диагностика: cargo +nightly rustc -- -Zprint-type-sizes, лечение — Box::pin в середине цепочки. Известная застарелая проблема компилятора: слоты в автомате переиспользуются далеко не идеально.

8. Неограниченная конкурентность и голодание

Async делает создание задач настолько дешёвым, что ограничения приходится ставить руками. Классика — for url in urls { tokio::spawn(fetch(url)) } на списке из миллиона URL: миллион задач в очереди, миллион TCP-соединений, EMFILE, OOM-killer.

use futures::stream::{self, StreamExt};

// Конкурентность ограничена явно: не более 32 запросов «в воздухе».
let bodies: Vec<_> = stream::iter(urls)
    .map(|url| async move { reqwest::get(url).await?.text().await })
    .buffer_unordered(32)
    .collect()
    .await;

Тот же приём другими средствами — tokio::sync::Semaphore с фиксированным числом разрешений и ограниченные каналы mpsc::channel(n) вместо unbounded_channel. Ограниченный канал даёт обратное давление бесплатно: производитель ждёт на send().await, когда потребитель не успевает. Неограниченный превращает отставание потребителя в рост потребления памяти — та же ошибка, что и в очередях сообщений (трек distributed-systems).

Обратная сторона — голодание. У tokio есть бюджет кооперации: 128 операций ввода-вывода на один poll задачи, после чего задача принудительно возвращает Pending, чтобы дать шанс соседям. Бюджет спасает от очевидных случаев, но не от вашего собственного цикла:

loop {
    // Если очередь никогда не пустеет, эта ветка всегда готова,
    // а всё остальное в этой задаче не исполняется никогда.
    if let Some(job) = queue.pop() {
        process(job);                       // синхронно и быстро — но всегда
        tokio::task::yield_now().await;     // ...поэтому явно уступаем воркер
    } else {
        notify.notified().await;
    }
}

Ещё одна тонкость: select! выбирает готовую ветку в случайном порядке именно чтобы избежать голодания. Если вам нужен детерминированный приоритет — добавьте biased; первой строкой, но тогда голодание младших ветвей становится вашей ответственностью.

Галерея подводных камней

Диагностика: чем смотреть внутрь

Обычный профилировщик покажет вам стек рабочего потока tokio, а не логическую задачу, — это принципиальное неудобство async. Инструменты, которые действительно помогают:

  • tokio-console — top для задач: список, время в состоянии Idle и Busy, количество опросов, предупреждения вида «задача не уступала воркер 500 мс» и «задача разбужена собственным wake без прогресса». Включается через console-subscriber и сборку с RUSTFLAGS="--cfg tokio_unstable". Первое, что стоит подключить при жалобах на задержки.
  • tracing с #[instrument] — единственный вменяемый способ получить связную картину: обычный лог из async перемешан между задачами, а span привязывает записи к конкретной операции сквозь все await.
  • tokio::time::timeout вокруг любого сетевого вызова — не диагностика, а гигиена: без него зависший запрос держит задачу и её ресурсы вечно.
  • Метрики рантайма (Handle::metrics): глубина очередей, число парковок, украденных задач. Плюс общие приёмы измерения из «Измерения» и наблюдаемость на уровне ОС из «Observability и производительность».

Симптом «p50 отличный, p99 в сто раз хуже» в асинхронном сервисе почти всегда означает блокировку воркера, а не медленный код.

Когда async не нужен

Async — это не «современный способ писать Rust». Это инструмент под конкретную форму нагрузки: много одновременных ожиданий.

Практические ориентиры:

  • CLI-утилита, которая делает пять HTTP-запросов — не нужен async. std::thread::spawn на каждый или просто последовательно. Вы сэкономите минуту работы и тридцать секунд времени сборки.
  • Сервер на 200 запросов в секунду с базой — можно и на потоках. Узкое место — база, не переключение контекста. Но здесь async уже не мешает, а экосистема (axum, sqlx) асинхронная, так что выбор скорее по библиотекам.
  • Числодробилка, обработка изображений, компиляция — нужен rayon, а не async. Async не ускоряет вычисления.
  • Прокси, шлюз, брокер, WebSocket-хаб, краулер — async, и без вариантов.
  • Микроконтроллерembassy даёт кооперативную многозадачность вообще без ОС и без аллокатора, и это один из самых чистых примеров пользы async: одна задача на датчик вместо ручного автомата состояний в обработчике прерывания (трек embedded, обзор экосистемы Rust).

Что это стоит: честно

Что вы получаете: миллион одновременных ожиданий на обычном сервере; отсутствие гонок данных даже в асинхронном коде (те же Send/Sync, что и для потоков, — подробно в предыдущей статье); отсутствие пауз сборщика мусора и предсказуемая задержка; лучшие в отрасли результаты по памяти на соединение; работоспособность вплоть до no_std.

Что вы платите:

  • Раскраска функций. async заразен снизу вверх: одна асинхронная функция в глубине требует async от всех вызывающих. Обратно ходить нельзя без block_on, а block_on внутри рантайма паникует. Развёрнутая критика — «Async Rust Is A Bad Language»; ответ на неё с проектной стороны — Withoutboats, «Why Async Rust». Прочитайте оба, они честнее любого учебника.
  • Сложность сообщений компилятора. Ошибки про Send в сгенерированной future — самые тяжёлые для чтения в языке: они указывают на строку tokio::spawn, а виновата другая строка. Навык чтения диагностики здесь окупается больше всего.
  • Время сборки. Пустой axum-сервис тянет 200+ крейтов и собирается с нуля минуты. Мономорфизация асинхронных цепочек генерирует много кода. Меры те же: воркспейс, cargo check, линкер mold, sccache.
  • Незакрытые дыры. Нет асинхронного Drop. Stream до сих пор не в std (используется futures::Stream и tokio_stream). async fn в трейтах не выражает Send. Async-блоки в замыканиях требуют приседаний. Это не мелочи, а то, обо что вы стукнётесь на второй неделе.
  • Отладка. Стек показывает воркер, а не логическую задачу. Без tracing и tokio-console вы почти слепы.

Куда async тащить не стоит: одноразовые скрипты; вычислительные конвейеры; программы с единицами одновременных операций; библиотеки общего назначения, которые иначе окажутся привязаны к конкретному рантайму (для них правильнее синхронный API, а асинхронность оставить приложению). Проектная эвристика: начните с потоков; переходите на async, когда упёрлись в число одновременных ожиданий, а не потому, что так модно. Общий взгляд на выбор модели конкурентности — в «Потоки и синхронизация».

Мини-итог

  • Future — трейт с одним методом poll, возвращающим Ready или Pending. Future инертна: без внешнего poll не происходит ничего.
  • async fn компилируется в автомат состояний. Переменные, живущие через .await, лежат внутри объекта future; не живущие — в обычном кадре стека. Внутренние ссылки делают future самоссылающейся, поэтому после первого poll её нельзя двигать: отсюда Pin, Box::pin, pin!.
  • Waker — единственный канал связи «я готов»: реактор хранит его и дёргает по событию из epoll или колеса таймеров. Исполнителя в std нет, рантайм — библиотека; на практике это tokio, и выбирается он один раз на проект.
  • Воркеров столько, сколько ядер, и они не вытесняются. Точка переключения ровно одна — .await.
  • Блокирующий вызов внутри async — ошибка номер один. Диагноз: p99 в разы хуже p50 под нагрузкой. Лечение: spawn_blocking.
  • MutexGuard из std через await ломает Send. Сначала пытайтесь сузить критическую секцию, tokio::sync::Mutex — второй выбор.
  • Отмена — это Drop в точке приостановки. Проверяйте cancellation safety всего, что кладёте в select!. Уронить JoinHandle — не значит отменить задачу.
  • join! — конкурентность в одной задаче; spawn — ещё и параллелизм, ценой Send + 'static. JoinSet отменяет свои задачи при уронении.
  • Ограничивайте конкурентность явно: buffer_unordered, Semaphore, ограниченные каналы. Обратное давление сама себе программа не создаст.
  • Async нужен, когда одновременных ожиданий много. Для вычислений — rayon, для единиц операций — потоки.

Источники

Что дальше

unsafe и FFI: когда необходимо, как ограничить, связь с C — спустимся ещё на уровень ниже. Мы уже видели, что Pin и самоссылающиеся future живут на границе того, что система типов умеет доказывать; дальше речь пойдёт о том, что делать, когда доказать нельзя вообще: как выглядит контракт unsafe, как обернуть сырой указатель в безопасный API, как позвать библиотеку на C и не потерять при этом гарантии, за которые вы платили всю предыдущую книгу.

Нашли неточность? Выделите фрагмент текста — рядом появится жучок.

Нужен разбор именно вашей ситуации?

Статья описывает общий случай. Если у вас частный — можно разобрать его отдельно, платно. А если не хватает целого материала, предложите тему: её оплачивают вскладчину, и она выходит открытой для всех.

Доска запросов