Каналы: паттерны и подводные камни

Каналы: паттерны и подводные камни

Каналы - основной способ коммуникации между горутинами. Но с ними легко выстрелить себе в ногу.

Buffered vs Unbuffered

// Unbuffered: отправитель блокируется, пока получатель не прочитает
ch := make(chan int)

// Buffered: отправитель блокируется только когда буфер полон
ch := make(chan int, 10)

Правило: unbuffered для синхронизации, buffered для очередей задач.

<?php

declare(strict_types=1);

namespace App\Concurrency;

use Symfony\Component\Messenger\MessageBusInterface;

final readonly class TaskMessage
{
    public function __construct(public int $value) {}
}

final readonly class ChannelLikeProducer
{
    public function __construct(private MessageBusInterface $bus) {}

    public function send(int $value): void
    {
        // эквивалент ch <- value - кладёт в очередь Redis/AMQP
        $this->bus->dispatch(new TaskMessage($value));
    }
}

В PHP нет встроенных каналов как примитива межгорутинного общения. Их роль играют брокеры очередей через Symfony Messenger transport: Redis Streams, RabbitMQ, AMQP, Doctrine (БД). «Buffered» канал = очередь с persistence, «unbuffered» = sync transport (выполнение в текущем процессе). Это межпроцессное общение, а не in-process - latency измеряется в миллисекундах, а не наносекундах как в Go.

Конфигурация transport (config/packages/messenger.yaml):

framework:
    messenger:
        transports:
            async:
                dsn: '%env(MESSENGER_TRANSPORT_DSN)%' # redis://redis:6379/messages
                options:
                    stream: 'tasks'
                    consumer: 'workers'
        routing:
            'App\Concurrency\TaskMessage': async

Unbuffered = рукопожатие sender-receiver; buffered = очередь до N сообщений

Directional channels

// Только для отправки
func producer(out chan<- int) {
    for i := 0; i < 10; i++ {
        out <- i
    }
    close(out)
}

// Только для чтения
func consumer(in <-chan int) {
    for val := range in {
        fmt.Println(val)
    }
}

Select: мультиплексирование каналов

select блокируется до любого готового case; если готовы несколько - выбирает случайный

func main() {
    ch1 := make(chan string)
    ch2 := make(chan string)

    go func() {
        time.Sleep(1 * time.Second)
        ch1 <- "one"
    }()

    go func() {
        time.Sleep(2 * time.Second)
        ch2 <- "two"
    }()

    for i := 0; i < 2; i++ {
        select {
        case msg := <-ch1:
            fmt.Println("ch1:", msg)
        case msg := <-ch2:
            fmt.Println("ch2:", msg)
        }
    }
}
<?php
declare(strict_types=1);

// PHP-эквивалент select - Promise\any() в ReactPHP: резолвится первым из списка.
use function React\Promise\any;
use React\EventLoop\Loop;
use React\Promise\Deferred;
use React\Promise\PromiseInterface;

final class Multiplexer
{
    /** Возвращает результат первого resolved promise (как select из двух каналов). */
    public function firstOf(PromiseInterface $a, PromiseInterface $b): PromiseInterface
    {
        return any([$a, $b]);
    }

    public function example(): void
    {
        $ch1 = new Deferred();
        Loop::addTimer(1.0, static fn () => $ch1->resolve('one'));

        $ch2 = new Deferred();
        Loop::addTimer(2.0, static fn () => $ch2->resolve('two'));

        $this->firstOf($ch1->promise(), $ch2->promise())
            ->then(static fn (string $msg) => print('first: ' . $msg . "\n"));
    }
}

Таймауты через select

func fetchWithTimeout(url string, timeout time.Duration) (string, error) {
    ch := make(chan string, 1)
    errCh := make(chan error, 1)

    go func() {
        result, err := fetch(url)
        if err != nil {
            errCh <- err
            return
        }
        ch <- result
    }()

    select {
    case result := <-ch:
        return result, nil
    case err := <-errCh:
        return "", err
    case <-time.After(timeout):
        return "", fmt.Errorf("timeout after %v", timeout)
    }
}
<?php

declare(strict_types=1);

namespace App\Concurrency;

use Symfony\Component\HttpClient\Exception\TimeoutException;
use Symfony\Contracts\HttpClient\HttpClientInterface;

final readonly class FetchWithTimeout
{
    public function __construct(private HttpClientInterface $client) {}

    public function fetch(string $url, float $timeoutSec): string
    {
        try {
            $response = $this->client->request('GET', $url, [
                'timeout' => $timeoutSec,
            ]);

            return $response->getContent();
        } catch (TimeoutException $e) {
            throw new \RuntimeException(\sprintf('timeout after %.2fs', $timeoutSec), 0, $e);
        }
    }
}

Прямого аналога select для каналов в PHP нет. На уровне HTTP-клиента таймаут задаётся параметром (Symfony HttpClient timeout). Если нужно «гонка нескольких Promise» - ReactPHP Promise\race или Promise\any. Symfony HttpClient по умолчанию sync, но поддерживает batch через stream() - близкий аналог fan-in.

ReactPHP-вариант с гонкой:

<?php

declare(strict_types=1);

namespace App\Concurrency;

use React\EventLoop\Loop;
use React\Promise\Deferred;
use React\Promise\PromiseInterface;

use function React\Promise\race;

final readonly class FetchWithRace
{
    public function fetch(PromiseInterface $work, float $timeoutSec): PromiseInterface
    {
        $timeout = new Deferred();
        Loop::addTimer($timeoutSec, static fn () => $timeout->reject(new \RuntimeException('timeout')));

        return race([$work, $timeout->promise()]);
    }
}

Done-channel паттерн

func worker(done <-chan struct{}) {
    for {
        select {
        case <-done:
            fmt.Println("shutting down")
            return
        default:
            // работа
            time.Sleep(100 * time.Millisecond)
        }
    }
}

func main() {
    done := make(chan struct{})

    go worker(done)

    time.Sleep(time.Second)
    close(done) // сигнал всем горутинам
}
<?php

declare(strict_types=1);

namespace App\Concurrency;

final class GracefulWorker
{
    private bool $shouldStop = false;

    public function __construct()
    {
        \pcntl_async_signals(true);
        \pcntl_signal(\SIGTERM, fn () => $this->shouldStop = true);
        \pcntl_signal(\SIGINT, fn () => $this->shouldStop = true);
    }

    public function run(): void
    {
        while (!$this->shouldStop) {
            $this->doWork();
            \usleep(100_000);
        }
    }

    private function doWork(): void
    {
        // тяжёлая итерация
    }
}

В PHP done-канал ближе всего к сигналу OS (SIGTERM/SIGINT) для воркера. Symfony Messenger обрабатывает их штатно: messenger:consume поллит флаг shouldStop() после каждого сообщения. Внутри воркера для долгих задач используют pcntl_signal + проверку флага.

Классический deadlock: отправляешь в небуферизированный канал в той же горутине, где читаешь. Go runtime его обнаружит и упадёт с `fatal error: all goroutines are asleep`.

Nil channel

// Отправка и чтение из nil-канала блокируются навсегда
// Это полезно для отключения ветки в select
var ch chan int // nil

select {
case v := <-ch: // никогда не сработает
    fmt.Println(v)
case <-time.After(time.Second):
    fmt.Println("timeout")
}
<?php
declare(strict_types=1);

// В PHP «nil-канала» нет. Эквивалентом «отключения ветки» в Promise\any
// служит передача только активных promises - кандидаты на race собираются
// в список conditionally.
use function React\Promise\any;
use React\EventLoop\Loop;
use React\Promise\Deferred;
use React\Promise\PromiseInterface;

final class ConditionalRace
{
    /**
     * @param list<PromiseInterface> $candidates  только те ветки, которые мы хотим
     *                                            оставить активными (как ненулевые каналы)
     */
    public function raceActive(array $candidates, float $timeoutSec): PromiseInterface
    {
        $timeout = new Deferred();
        Loop::addTimer($timeoutSec, static fn () => $timeout->reject(new \RuntimeException('timeout')));

        return any([...$candidates, $timeout->promise()]);
    }
}

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