Сервис читал события из очередей вроде Kafka, Redis Streams или NATS, и для каждого события запускал асинхронную задачу в Tokio. Внутри такой задачи происходил fan-out: до 1000 пользовательских токенов — под каждый создавалась отдельная Tokio-задача, которая делала внешний вызов и ждала ответа. Потом все ответы собирались через JoinSet, и формировалось выходное событие. Код выглядел безобидно: токен-задачи короткоживущие, ответы приходят за миллисекунды, и казалось, что ранние события должны целиком завершаться раньше, даже при лавине новых.
Однако логи показали другую картину. При всплеске из 1000 событий суммарно порождалось около миллиона задач. Задачи от первого события начинали выполняться далеко не первыми — часть из них стартовала только ближе к концу всей пачки. При этом всплеск укладывался в отведённое время, и throughput не страдал, но пиковое потребление памяти оказалось выше ожидаемого.
Причина — в устройстве планировщика Tokio. Многопоточный рантайм использует локальные очереди воркеров (до 256 задач каждая) и общую глобальную очередь. Когда локальная очередь переполняется, половина задач уходит в глобальную, откуда их могут подхватить любые воркеры. Плюс работает work stealing. Tokio не знает, из какого события пришла задача: для него все токен-задачи, задачи родительских событий, ждущие JoinSet, и просыпающиеся по I/O — просто runnable tasks, конкурирующие за poll. Поэтому создание задачи вовсе не означает её скорого первого опроса, а первый опрос — скорого завершения.
Из-за такого перемешивания некоторые токен-задачи из ранних событий задерживались до самого конца всплеска. Пока они висели, их родительские событийные задачи тоже оставались живы, удерживая в памяти всё состояние события. Тысячи одновременно живых задач, пусть и с небольшим состоянием каждая, в сумме давали заметный скачок peak memory.
Справедливое планирование Tokio гарантирует только при ограниченном количестве задач и отсутствии блокировок. Ограничивать должен сам разработчик. Автор добавил Semaphore, разрешающий одновременную обработку только N событий. Теперь новые события не поступали в обработку, пока старые не завершатся, и все токен-задачи одного события завершались близко друг к другу. Подбор числа N делался экспериментально. Важно, что с этим ограничением пиковая память сильно снизилась, а всплеск по-прежнему укладывался в требуемое время. Опасения, что Semaphore зарежет пропускную способность, не подтвердились.
Вывод прост: запустить задачу рано — не значит получить её первый poll рано. А первый poll рано — не значит раннее завершение. Если в приложении есть своя единица справедливости (в данном случае событие), Tokio о ней не знает, и границы нужно выставлять самому. И перед каждым spawn стоит прикинуть, сколько задач может висеть одновременно, потому что это напрямую влияет на память, а иногда и на другие удерживаемые ресурсы.