Tokio гарантирует продвижение, а не порядок
- title
- Tokio гарантирует продвижение, а не порядок
- type
- summary
- summary
- Гарантия честности в Tokio предполагает ограничение числа задач, которое может задать только приложение
- tags
- rust, async, concurrency, performance
- sources
- tokio-progress-not-ordering
- created
- 2026-07-29
- updated
- 2026-07-29
- lang
- ru
- translation_of
- tokio-progress-not-ordering
- source_updated
- 2026-07-29
- translated
- 2026-09-01
- translator
- lllm/antigravity/gemini-3.7-flash-medium
Рассказ Пранитхи о резком скачке памяти в event-driven сервисе на Rust. Заметка написана как предыстория к её статье о поведении аллокатора - это расследование со стороны приложения, проведённое ещё до того, как выяснилось, что память связана с аллокатором.
Сервис вычитывал события из очереди (Kafka, Redis Streams или NATS) и обрабатывал каждое с веерным ветвлением (fan-out). Каждое событие содержало полезную нагрузку около 4 КБ и до 1000 токенов пользователей; для каждого токена требовался исходящий вызов, а ответы затем собирались обратно:
loop {
let event: Event = fetch_next_event().await;
tokio::spawn(async move {
let mut tasks = JoinSet::new();
for token in &event.user_tokens {
tasks.spawn(async move { process(token, data).await });
}
let mut responses = Vec::with_capacity(event.user_tokens.len());
while let Some(res) = tasks.join_next().await {
responses.push(res)
}
generate_response_event(event, responses);
});
}
Веерное ветвление ничем не ограничено, но автор полагала, что на практике это не создаст проблем. Задачи живут недолго - несколько миллисекунд каждая, после чего завершаются, - поэтому задачи, запущенные раньше, должны и закончиться раньше. Завершение отдельных исходящих вызовов не по порядку ожидалось; предположение, что более ранние события завершатся раньше поздних, казалось надёжным.
Что показали логи
Во время всплеска из 1000 событий примерно по 1000 токенов в каждом - около 1 млн задач - записи в логах перемешались следующим образом:
finished: event 779
finished: event 976
started: event 900, user 42
started: event 900, user 261
started: event 1, user 974
started: event 1, user 831
Задачи для токенов из события 1 получили свой первый poll уже после того, как события 779 и 976 полностью завершились. Всплеск всё равно уложился в отведённый бюджет времени, так что проблема была не в задержках. Неожиданным оказался именно разрыв между порядком отправки задач и их первым опросом (poll).
Почему планировщик ведёт себя именно так
Многопоточный runtime Tokio имеет фиксированный набор worker-потоков, локальную очередь на каждый worker ёмкостью 256 задач и общую глобальную очередь. При переполнении локальной очереди worker перемещает половину своих задач в глобальную очередь. Потоки-worker'ы отдают предпочтение собственной локальной очереди, периодически проверяют глобальную и перехватывают (steal) задачи у других worker'ов во время простоя.
Как только задача становится независимой единицей планирования, runtime уже понятия не имеет, какое событие её породило. 1000 задач для токенов из одного события перемешиваются с задачами для токенов из всех остальных событий, с родительскими задачами событий, ждущими на JoinSet, и с задачами, просыпающимися по готовности I/O, - всё это просто готовая к исполнению работа, конкурирующая за место в очереди. Переполнение очередей и work stealing меняют порядок взятия задач в работу относительно порядка их отправки. Вывод, к которому приходит автор:
task created != task polled != task completed
Следствие для памяти очевидно. Каждая задача несёт в себе состояние, само по себе небольшое, а пиковое потребление памяти зависит от того, сколько задач существует одновременно, а не от того, сколько из них выполняется прямо сейчас. Несколько задач для токенов из ранних событий, доживших до конца всплеска, не давали завершиться родительским задачам событий, а те, в свою очередь, удерживали полезную нагрузку в 4 КБ и векторы токенов.
У гарантии честности есть предварительное условие
Задокументированная честность планировщика Tokio работает при ограниченном количестве задач и при условии, что ни одна задача не блокирует поток-worker. В исходном коде ограничений не было вовсе: события вычитывались с максимальной скоростью, каждое порождало задачу, а та - ещё до 1000 новых. Tokio принимает всё, что ему дают, и обеспечивает продвижение всей этой массы задач; ограничение объёма, на которое опирается гарантия, должно исходить от самого приложения.
На самом же деле требовалась честность на уровне событий: чтобы все задачи токенов, принадлежащие одному событию, опрашивались и завершались примерно в то время, когда событие поступило. Этого удалось добиться с помощью Semaphore, ограничивающего число одновременно обрабатываемых событий. Количество разрешений (permits) подобрали опытным путём, поскольку подходящее число зависит от того, сколько времени занимает возврат из poll у каждой задачи. Пропускная способность не пострадала: всплеск завершился в срок, а пиковое потребление памяти заметно снизилось.
Почему это легко упустить
Никаких явных ошибок не было. Ни одна задача не потерялась, ни один запрос не тормозил, целевая пропускная способность выдерживалась. Проблема проявилась исключительно как скачок памяти без видимых причин в логике самого приложения: возникло несоответствие между ментальной моделью автора (событие - это единица работы, которая начинается и заканчивается) и моделью runtime'а (задача есть задача).
Такова общая структура проблемы: если на уровне приложения есть единица честности, соблюдения которой вы ждёте от runtime'а, то сам runtime о её существовании даже не догадывается. В safety-in-an-unsafe-world утверждается, что библиотеки должны кодировать свои инварианты в типах, чтобы нарушение приводило к ошибке компиляции; здесь же обратный случай: инвариант принадлежит вызывающей стороне, и Tokio не может предоставить тип для его выражения. На практике перед каждым tokio::spawn стоит задаваться вопросом: каково максимальное число одновременно существующих задач, которое здесь может возникнуть, и что удерживает в памяти каждая из них во время ожидания?