Асинхронные вызовы
Winter переносит в PHP модель конкурентности из Java: ExecutorService, Future,
CompletableFuture и аннотацию @Async. Имена, сигнатуры и семантика взяты у
java.util.concurrent намеренно — это проверенный десятилетиями словарь, который
не нужно изобретать заново.
Что это и зачем
Асинхронный вызов — вызов метода, который возвращает управление, не дожидаясь окончания работы.
Проблема. В обработчике почти всегда есть работа, ради которой клиент не должен ждать. Пользователь зарегистрировался — надо создать запись, отправить письмо, дёрнуть CRM, положить событие в аналитику. Значение имеет только первое; остальное добавляет к ответу секунды, а при недоступности внешнего сервиса ещё и роняет запрос, который по сути был успешным.
Другая сторона той же задачи — обратная. Странице нужны данные из трёх сервисов, и они друг от друга не зависят. Последовательные вызовы дают сумму трёх задержек там, где достаточно самой большой из них.
Решение. Пометьте метод атрибутом — и вызывающий код продолжит работу сразу.
Когда результат всё-таки нужен, метод возвращает Future, у которого этот результат
можно спросить позже.
Родословная
Если вы писали на Java, этот раздел вам уже знаком — совпадают не только идеи, но и названия:
| Winter | Аналог в Java |
|---|---|
#[Async] |
@Async из Spring |
#[EnableAsync] |
@EnableAsync из Spring |
Future |
java.util.concurrent.Future |
CompletableFuture |
java.util.concurrent.CompletableFuture |
ExecutorService |
java.util.concurrent.ExecutorService |
Executors |
java.util.concurrent.Executors |
RejectPolicy |
Политики ThreadPoolExecutor — AbortPolicy, CallerRunsPolicy, DiscardPolicy |
ExecutionException, TimeoutException, CancellationException, RejectedExecutionException |
Одноимённые из java.util.concurrent |
Сигнатуры тоже совпадают: future.get(timeout), future.cancel(mayInterruptIfRunning),
executor.submit(...), executor.invokeAll(...), CompletableFuture.allOf(...).
Где Winter отличается от Java
Три отличия, о которых стоит знать заранее.
Корутины вместо потоков. Задача выполняется не в отдельном потоке ОС, а в корутине Swoole. Это кооперативная многозадачность: две задачи никогда не исполняют код одновременно, они чередуются, когда одна упирается в ожидание — база, сеть, файл, пауза. Отсюда правило: асинхронность ускоряет ожидание, но не вычисления. Два обращения к внешним API пойдут параллельно; два тяжёлых расчёта — нет.
Композиции нет. thenApply, thenCompose и остальной конвейер стадий не
перенесены сознательно: в PHP они выродились бы в лапшу из замыканий. Есть
whenComplete() для колбэка и allOf() для ожидания набора — этого хватает
подавляющему большинству задач.
Вызов соседнего метода работает. В Spring @Async не срабатывает, когда метод
зовут через this внутри того же бина, — там прокси-обёртка. Winter подменяет класс
наследником, объект остаётся один, поэтому $this->asyncMethod() остаётся
асинхронным. Единственное условие — метод не должен быть private.
Это не замена очереди
Работа остаётся внутри того же воркера и живёт ровно столько, сколько он. Перезапуск её потеряет, на другую машину она не уедет, повторить её после сбоя некому.
Для работы, которая обязана быть выполненной, нужна очередь и
процесс или демон. #[Async] —
про то, чтобы не заставлять ждать, а не про гарантии доставки.
Быстрый старт
Два шага: включить проксирование и пометить метод.
use Flytachi\Winter\Kernel\App\Attribute\EnableAsync;
#[EnableWeb]
#[EnableAsync]
final class Application extends WinterApplication { /* ... */ }<?php
namespace Main;
use Flytachi\Winter\DI\Attribute\Autowired;
use Flytachi\Winter\Kernel\Concurrent\Async\Async;
class NotificationService
{
#[Autowired] private Mailer $mailer;
#[Async]
public function sendWelcome(User $user): void
{
$this->mailer->send($user->email, 'welcome');
}
}Вызов ничем не отличается от обычного — этим и удобен:
#[PostMapping]
public function register(#[RequestJson, Valid] RegisterDto $dto): ResponseEntity
{
$user = $this->users->create($dto);
$this->notifications->sendWelcome($user); // не задержит ответ
return ResponseEntity::created($user);
}Без `#[EnableAsync]` атрибут молчит
Как и @EnableAsync в Spring, этот атрибут включает механизм целиком. Без него
методы с #[Async] выполняются синхронно, и никакой ошибки не будет —
приложение просто работает медленнее, чем вы думаете.
«Потом» иногда значит «прямо сейчас»
Swoole переключается в новую корутину немедленно. Если тело метода не содержит ни одного ожидания — нет обращений к сети, базе или паузы, — оно успеет полностью отработать до того, как вызов вернёт управление. Результат от этого верен, но интуиция «выполнится позже» не работает.
Результат работы
Метод объявляет одно из двух — и ничего третьего.
void — отправил и забыл
Результат недоступен, исключение внутри уйдёт только в лог.
#[Async]
public function trackEvent(string $name, array $payload): void
{
$this->analytics->push($name, $payload);
}Годится для всего, чей исход не влияет на ответ: метрики, уведомления, прогрев кеша.
Future — результат понадобится позже
use Flytachi\Winter\Kernel\Concurrent\{CompletableFuture, Future};
#[Async]
public function fetchProfile(int $id): Future
{
return CompletableFuture::completedFuture($this->api->profile($id));
}Это и даёт вторую форму выигрыша — запустить независимые операции разом:
$profile = $this->users->fetchProfile($id); // пошло
$orders = $this->orders->fetchRecent($id); // пошло
$balance = $this->billing->fetchBalance($id); // пошло
return ResponseEntity::ok([ // ждём здесь
'profile' => $profile->get(),
'orders' => $orders->get(),
'balance' => $balance->get(),
]);Три обращения к внешним сервисам идут одновременно, и запрос занимает время самого медленного, а не сумму трёх.
Тип возврата проверяется строго
Допустимы ровно void и Future. Ни ?Future, ни CompletableFuture, ни array
не подойдут: прокси не сгенерируется, и приложение упадёт при загрузке с
AsyncException. Это лучше тихой поломки — ошибку видно сразу.
Future и CompletableFuture
Future — то же обещание, что в Java: контракт на значение, которого пока нет.
| Метод | Что делает |
|---|---|
get(?float $timeout = null) |
Ждёт и возвращает результат |
isDone() |
Завершилась ли — успешно, с ошибкой или отменой |
isCancelled() |
Была ли отменена |
cancel(bool $mayInterruptIfRunning = false) |
Пытается отменить |
get() — единственный способ получить значение, и он же место, где всплывают
ошибки:
| Исключение | Когда |
|---|---|
ExecutionException |
Задача завершилась броском; оригинал — в getPrevious() |
TimeoutException |
Истёк переданный таймаут; задача при этом продолжает выполняться |
CancellationException |
Задачу отменили |
try {
$data = $future->get(timeout: 2.0);
} catch (TimeoutException) {
$data = $this->cache->stale($id); // не дождались — отдаём что есть
} catch (ExecutionException $e) {
$this->logger->error('fetch failed', ['cause' => $e->getPrevious()?->getMessage()]);
$data = null;
}CompletableFuture — реализация, которой можно управлять вручную:
| Метод | Что делает |
|---|---|
completedFuture($value) |
Готовое обещание с результатом |
failedFuture($throwable) |
Готовое обещание с ошибкой |
supplyAsync($fn, $executor = null) |
Запустить и вернуть результат |
runAsync($fn, $executor = null) |
Запустить, результат отбросить |
allOf(Future ...$futures) |
Обещание, завершающееся, когда завершились все |
join() |
То же, что get() без таймаута |
whenComplete($fn) |
Колбэк, получающий значение и ошибку |
complete($value) / completeExceptionally($e) |
Завершить вручную |
Дождаться набора:
CompletableFuture::allOf($profile, $orders, $balance)->get(timeout: 5.0);Таймаут здесь общий на весь набор.
Пулы исполнителей
По умолчанию задачи уходят в общий исполнитель. Когда нужен контроль — ограничить число одновременных обращений к внешнему API, отделить один вид работы от другого, — заведите собственный пул и укажите его имя в атрибуте.
use Flytachi\Winter\Kernel\App\Attribute\{Bean, Configuration};
use Flytachi\Winter\Kernel\Concurrent\{ExecutorService, Executors, RejectPolicy};
#[Configuration]
final class AppConfig
{
#[Bean(name: 'pool.sms')]
public function smsPool(): ExecutorService
{
return Executors::newFixedExecutor(
concurrency: 5, // не больше 5 одновременно
queue: 100, // и ещё 100 в ожидании
onReject: RejectPolicy::CALLER_RUNS,
);
}
}#[Async('pool.sms')]
public function sendSms(string $phone, string $text): void { /* ... */ }Имя в атрибуте — это ключ в контейнере, поэтому пул может быть любым объектом,
реализующим ExecutorService.
| Фабрика | Что даёт |
|---|---|
Executors::common() |
Общий исполнитель, без ограничений |
Executors::newFixedExecutor($concurrency, $queue, $onReject) |
Пул с пределом параллелизма и очередью |
Executors::newCoroutineExecutor() |
Отдельный корутинный исполнитель |
Executors::shutdownCommon($timeout) |
Остановить общий исполнитель |
Есть и Executors::newDeferredExecutor() — исполнитель, выполняющий задачи после
отправки ответа. Он относится к запуску под PHP-FPM и в Swoole-приложении не нужен.
Пул с ограничением умеет рассказать о себе — удобно для метрик:
| Метод | Что возвращает |
|---|---|
concurrency() |
Предел одновременных задач |
activeCount() |
Сколько выполняется сейчас |
queuedCount() |
Сколько принято и ждёт своей очереди |
remainingCapacity() |
Сколько ещё примет |
Что делать при переполнении
Когда заполнены и параллелизм, и очередь, срабатывает политика отказа — те же три,
что у ThreadPoolExecutor в Java:
| Политика | Поведение |
|---|---|
ABORT (по умолчанию) |
Бросает RejectedExecutionException |
CALLER_RUNS |
Выполняет задачу синхронно у вызывающего — естественное торможение |
DISCARD |
Молча отбрасывает; вернётся отменённая Future |
CALLER_RUNS обычно самый здравый выбор для веба: при перегрузке эндпоинт станет
медленнее, но не начнёт терять работу и не будет отдавать ошибки.
Очередь по умолчанию не ограничена
queue: 0 означает, что задачи копятся без предела и отказа не наступит никогда.
Политика отказа имеет смысл только вместе с явным размером очереди.
Пул напрямую, без атрибута
#[Async] — удобная обёртка, но пул можно использовать и сам по себе: это тот же
ExecutorService, что в Java, с теми же методами.
| Метод | Что делает |
|---|---|
submit(callable $task, ...$args): Future |
Запускает и возвращает обещание результата |
execute(callable $task, ...$args): void |
Запускает, результат отбрасывает |
invokeAll(iterable $tasks, ?float $timeout = null): array |
Запускает набор и ждёт все; возвращает Future в порядке задач |
shutdown(): void |
Больше не принимать новые задачи |
isShutdown(): bool |
Остановлен ли |
awaitTermination(?float $timeout = null): bool |
Дождаться завершения принятых задач |
Это уместно там, где задач заранее неизвестное число и заводить под каждую метод незачем:
use Flytachi\Winter\Kernel\Concurrent\Executors;
$executor = Executors::common();
$futures = [];
foreach ($regions as $region) {
$futures[$region] = $executor->submit(fn() => $this->api->stats($region));
}
$stats = [];
foreach ($futures as $region => $future) {
$stats[$region] = $future->get(timeout: 5.0);
}invokeAll() делает то же самое короче, если результаты нужны все и сразу:
$futures = $executor->invokeAll(
array_map(fn($r) => fn() => $this->api->stats($r), $regions),
timeout: 5.0,
);Что с общим состоянием
Вопрос, который возникает первым у всех, кто писал многопоточный код: нужны ли блокировки?
Внутри воркера — нет. Корутины исполняются по очереди в одном потоке, поэтому
две задачи физически не могут изменять одну переменную одновременно. Гонок за память
в том смысле, в каком они бывают у потоков, здесь не существует, и synchronized
здесь нечему соответствовать.
Но переключение всё же есть. Оно происходит в точках ожидания — на обращении к базе, к сети, к файлу, на паузе. Значит, между двумя строками вашего кода, если между ними есть такое обращение, успеет отработать другая задача:
$balance = $this->repo->balance($id); // ← здесь другая задача может вклиниться
$this->repo->setBalance($id, $balance - $amount); // и прочитать старое значениеЭто не свойство #[Async], а свойство любой конкурентной обработки: то же самое
верно для двух одновременных HTTP-запросов. Лечится так же, как лечилось всегда, —
атомарной операцией в базе или блокировкой на её стороне, а не средствами PHP.
Состояние между воркерами не разделяется вовсе. Каждый воркер — отдельный процесс со своей памятью, поэтому статическое свойство, счётчик в объекте или локальный кеш видны только внутри него. Общее на всё приложение живёт снаружи: в базе, в Redis, в хранилище.
Область запроса задачу не наследует
Задача получает собственный корутинный контекст. #[Request]-объекты — контекст
аутентификации, текущая локаль — в неё не переходят: это сделано намеренно,
чтобы задача, пережившая ответ, не держала данные чужого запроса.
Всё нужное передавайте аргументами явно.
Требования к классу и методу
Подмена работает наследованием: фреймворк генерирует потомка вашего класса и переопределяет в нём помеченные методы. Отсюда все ограничения — они ровно те же, что у обычного наследования в PHP.
К классу:
| Нельзя | Сообщение при старте |
|---|---|
final класс |
the class is final and cannot be extended |
| Абстрактный класс, интерфейс, enum | only instantiable classes can be proxied |
| Встроенный класс PHP | internal classes cannot be proxied |
К методу:
| Нельзя | Почему |
|---|---|
final |
Потомок не сможет переопределить |
static |
Асинхронность применяется к экземпляру |
private |
Разрешается статически внутри своего класса — потомок его не перехватит. Сделайте protected: обращения к себе тоже пойдут через подмену |
abstract |
Помечайте реализацию |
Параметр по ссылке &$x |
Вызов возвращается раньше тела, записывать некуда |
К типу возврата — ровно два допустимых варианта, и оба обязаны быть указаны явно:
#[Async] public function send(): void { } // ✓ отправил и забыл
#[Async] public function fetch(): Future { } // ✓ результат заберут позже
#[Async] public function send() { } // ✗ типа возврата нет
#[Async] public function fetch(): ?Future { } // ✗ nullable не принимается
#[Async] public function count(): int { } // ✗ значение вернуть неоткудаОтсутствие типа — такая же ошибка, как неверный тип: the method has no return type. Привычка не писать : void здесь не сработает.
Нарушение видно при старте, а не в бою
Любое из этих ограничений роняет приложение на загрузке — с именем класса, именем метода, причиной и подсказкой, что сделать. В рантайме такая ошибка не всплывёт.
Проверить, не дожидаясь запуска, можно командой
call di build: она собирает те же подмены и падает с тем же
сообщением.
Тихая ловушка — new
Асинхронность обеспечивается подменой класса на потомка, а получить эту подмену можно только у контейнера.
// ✓ асинхронно — объект пришёл из контейнера
#[Autowired] private NotificationService $notifications;
// ✗ синхронно — обычный класс, атрибут проигнорирован
$notifications = new NotificationService();Коварство в том, что ничего не сломается: тот же тип, тот же результат, ни ошибки, ни предупреждения. Просто работа выполнится синхронно, и вызов будет ждать.
Фреймворк находит такие места статически:
php call di buildКоманда, помимо сборки прокси, показывает, где класс с #[Async] создаётся через
new. Проверка эвристическая — динамическое создание (new $class, фабрики) она не
видит, — но типовой случай ловит.
Примеры
Побочные эффекты после регистрации
<?php
namespace Main;
use Flytachi\Winter\DI\Attribute\Autowired;
use Flytachi\Winter\Kernel\Concurrent\Async\Async;
use Psr\Log\LoggerInterface;
class RegistrationService
{
#[Autowired] private Mailer $mailer;
#[Autowired] private CrmClient $crm;
#[Autowired] private LoggerInterface $logger;
public function register(RegisterDto $dto): User
{
$user = $this->users->create($dto);
$this->afterRegistration($user); // асинхронно, ответ не ждёт
return $user;
}
#[Async]
protected function afterRegistration(User $user): void
{
try {
$this->mailer->send($user->email, 'welcome');
$this->crm->createLead($user);
} catch (\Throwable $e) {
$this->logger->error('post-registration failed', [
'user' => $user->id,
'error' => $e->getMessage(),
]);
}
}
}Метод объявлен protected, а не private, — иначе подмена не сработала бы. И
try/catch внутри обязателен: у void-метода некому вернуть ошибку.
Параллельный сбор данных
<?php
namespace Main;
use Flytachi\Winter\DI\Attribute\Autowired;
use Flytachi\Winter\Kernel\Concurrent\Async\Async;
use Flytachi\Winter\Kernel\Concurrent\{CompletableFuture, Future, TimeoutException};
class DashboardService
{
#[Autowired] private BillingApi $billing;
#[Autowired] private StatsApi $stats;
#[Async]
public function balance(int $userId): Future
{
return CompletableFuture::completedFuture($this->billing->balance($userId));
}
#[Async]
public function usage(int $userId): Future
{
return CompletableFuture::completedFuture($this->stats->monthly($userId));
}
public function build(int $userId): array
{
$balance = $this->balance($userId);
$usage = $this->usage($userId);
try {
return [
'balance' => $balance->get(timeout: 3.0),
'usage' => $usage->get(timeout: 3.0),
];
} catch (TimeoutException) {
return ['balance' => null, 'usage' => null, 'degraded' => true];
}
}
}Два обращения идут одновременно, а таймаут не даёт странице зависнуть, если один из сервисов отвечает плохо.
Ограниченный доступ к внешнему шлюзу
Когда партнёрское API держит лимит по одновременным подключениям, превышать его нельзя, сколько бы запросов ни пришло.
#[Configuration]
final class GatewayConfig
{
#[Bean(name: 'pool.gateway')]
public function gatewayPool(): ExecutorService
{
return Executors::newFixedExecutor(
concurrency: 3,
queue: 50,
onReject: RejectPolicy::CALLER_RUNS,
);
}
}
class PaymentService
{
#[Async('pool.gateway')]
public function notifyGateway(Payment $payment): void { /* ... */ }
}Пул общий на воркер: сколько бы запросов его ни использовало, к шлюзу одновременно уйдёт не больше трёх. Переполнение очереди затормозит вызывающего, а не приведёт к отказу.
Ручное управление обещанием
CompletableFuture можно завершить самому — это удобно, когда результат приходит
не из вызова, а извне: из колбэка, из подписки, из другого механизма.
$future = new CompletableFuture();
$this->bus->subscribe('payment.confirmed', function (Payment $p) use ($future) {
$future->complete($p);
});
try {
$payment = $future->get(timeout: 30.0);
} catch (TimeoutException) {
throw new ResponseException('Payment confirmation timed out', HttpCode::GATEWAY_TIMEOUT);
}Вместе с планировщиком
#[Async] можно поставить на метод, уже помеченный #[Scheduled], — сочетание
рабочее. Но оно снимает главную гарантию планировщика: задача перестаёт быть
защищённой от наложения сама на себя, потому что вызов возвращается сразу и прогон
считается оконченным в момент отправки.
Разбор с правилами и выбором политики отказа — в разделе
Вместе с #[Async] на странице
планировщика.
Дальше
- Состав приложения — чем
#[Async]отличается от компонентов - Планировщик — совмещение с расписанием
- Процессы — когда работа обязана пережить запрос
- Демоны — когда её ещё и много
- Рантаймы — корутины Swoole, на которых всё это работает
- Внедрение зависимостей — почему объект нужно брать из контейнера