errgroup и graceful shutdown
errgroup и graceful shutdown
WaitGroup не умеет возвращать ошибки и отменять другие горутины при ошибке. Для этого есть errgroup поверх context.
errgroup: WaitGroup с ошибками
import "golang.org/x/sync/errgroup"
func fetchAll(ctx context.Context, urls []string) ([]string, error) {
g, ctx := errgroup.WithContext(ctx)
results := make([]string, len(urls))
for i, url := range urls {
i, url := i, url
g.Go(func() error {
body, err := fetchURL(ctx, url)
if err != nil {
return err // отменит ctx для всех
}
results[i] = body
return nil
})
}
if err := g.Wait(); err != nil {
return nil, err
}
return results, nil
}
<?php
declare(strict_types=1);
namespace App\Concurrency;
use React\Promise\PromiseInterface;
use Symfony\Contracts\HttpClient\HttpClientInterface;
use function React\Async\async;
use function React\Promise\all;
final readonly class ParallelFetcher
{
public function __construct(private HttpClientInterface $client) {}
/**
* @param list<string> $urls
* @return PromiseInterface<list<string>>
*/
public function fetchAll(array $urls): PromiseInterface
{
$promises = [];
foreach ($urls as $url) {
$promises[] = async(fn (): string => $this->client->request('GET', $url)->getContent())();
}
// all() резолвится при успехе всех или режектится при первой ошибке
return all($promises);
}
}
В PHP ближайший аналог - ReactPHP
Promise\all(ждёт все, fails-fast на первой ошибке) или Symfony HttpClientstream()с обработкой исключений. Cancellation в ReactPHP черезCancellablePromise. Symfony Messenger предлагаетMessageBatchдля группировки, но без авто-отмены - там скорее retry-логика, чем fail-fast.
Sync-вариант на Symfony HttpClient (без ReactPHP):
<?php
declare(strict_types=1);
namespace App\Concurrency;
use Symfony\Contracts\HttpClient\HttpClientInterface;
final readonly class SyncParallelFetcher
{
public function __construct(private HttpClientInterface $client) {}
/**
* @param list<string> $urls
* @return list<string>
*/
public function fetchAll(array $urls): array
{
$responses = [];
foreach ($urls as $i => $url) {
$responses[$i] = $this->client->request('GET', $url);
}
$results = [];
try {
foreach ($this->client->stream($responses) as $response => $chunk) {
if ($chunk->isLast()) {
$idx = \array_search($response, $responses, true);
$results[$idx] = $response->getContent();
}
}
} catch (\Throwable $e) {
// fail-fast: отменяем остальные
foreach ($responses as $r) {
$r->cancel();
}
throw $e;
}
return $results;
}
}
errgroup с лимитом конкурентности
g, ctx := errgroup.WithContext(ctx)
g.SetLimit(5) // максимум 5 горутин одновременно
for _, task := range tasks {
task := task
g.Go(func() error {
return processTask(ctx, task)
})
}
if err := g.Wait(); err != nil {
log.Fatal(err)
}
Graceful shutdown HTTP-сервера
func main() {
srv := &http.Server{Addr: ":8080"}
// Запускаем сервер в горутине
go func() {
if err := srv.ListenAndServe(); err != http.ErrServerClosed {
log.Fatalf("HTTP server error: %v", err)
}
}()
// Ждём сигнала завершения
ctx, stop := signal.NotifyContext(context.Background(),
syscall.SIGINT, syscall.SIGTERM,
)
defer stop()
<-ctx.Done()
log.Println("Shutting down...")
// Даём 10 секунд на завершение текущих запросов
shutdownCtx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
defer cancel()
if err := srv.Shutdown(shutdownCtx); err != nil {
log.Fatalf("Shutdown error: %v", err)
}
log.Println("Server stopped")
}
<?php
declare(strict_types=1);
namespace App\Concurrency;
use Psr\Http\Message\ServerRequestInterface;
use React\EventLoop\Loop;
use React\Http\HttpServer;
use React\Http\Message\Response;
use React\Socket\SocketServer;
final class GracefulHttpServer
{
public function start(string $listen, int $shutdownTimeoutSec): void
{
$http = new HttpServer(
static fn (ServerRequestInterface $request): Response => new Response(200, [], 'ok'),
);
$socket = new SocketServer($listen);
$http->listen($socket);
$shutdown = function () use ($socket, $shutdownTimeoutSec): void {
$socket->close(); // перестаём принимать соединения
Loop::addTimer(
$shutdownTimeoutSec,
static fn () => Loop::stop(),
);
};
Loop::addSignal(\SIGTERM, $shutdown);
Loop::addSignal(\SIGINT, $shutdown);
}
}
В классическом PHP-FPM graceful shutdown HTTP-сервера - это задача fpm и nginx, а не приложения:
php-fpmловит SIGQUIT, дожидается завершения запущенных запросов. Приложение участвует только если живёт долго - воркер Symfony Messenger или ReactPHP Http\Server. У Messenger флаги--time-limitи--memory-limit, у ReactPHP -Loop::addSignal.
Паттерн: запуск нескольких сервисов
func run(ctx context.Context) error {
g, ctx := errgroup.WithContext(ctx)
// HTTP сервер
g.Go(func() error {
return runHTTPServer(ctx, ":8080")
})
// gRPC сервер
g.Go(func() error {
return runGRPCServer(ctx, ":9090")
})
// Worker для очередей
g.Go(func() error {
return runConsumer(ctx)
})
return g.Wait()
}