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 HttpClient stream() с обработкой исключений. 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: первая ошибка отменяет общий ctx, оставшиеся горутины завершаются по ctx.Done(), Wait возвращает первую ошибку

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()
}
Любая горутина, которая может упасть - должна уметь сообщить об этом наверх. errgroup + context - стандартный способ.

Зарегистрируйтесь бесплатно, чтобы пройти квиз, решить задание с автопроверкой и вести прогресс.