Паттерны: fan-in, fan-out, pipeline, worker pool

Паттерны: fan-in, fan-out, pipeline, worker pool

Эти паттерны - строительные блоки для конкурентных систем поверх горутин и каналов.

Pipeline

Последовательная обработка данных через цепочку этапов.

Pipeline: generate → filter → multiply → consumer, каждая стадия в своей горутине, close распространяется вниз

// Генератор → Фильтр → Маппер → Потребитель

func generate(nums ...int) <-chan int {
    out := make(chan int)
    go func() {
        for _, n := range nums {
            out <- n
        }
        close(out)
    }()
    return out
}

func filter(in <-chan int, predicate func(int) bool) <-chan int {
    out := make(chan int)
    go func() {
        for n := range in {
            if predicate(n) {
                out <- n
            }
        }
        close(out)
    }()
    return out
}

func multiply(in <-chan int, factor int) <-chan int {
    out := make(chan int)
    go func() {
        for n := range in {
            out <- n * factor
        }
        close(out)
    }()
    return out
}

func main() {
    // Pipeline: генерируем → фильтруем чётные → умножаем на 10
    nums := generate(1, 2, 3, 4, 5, 6, 7, 8, 9, 10)
    evens := filter(nums, func(n int) bool { return n%2 == 0 })
    result := multiply(evens, 10)

    for v := range result {
        fmt.Println(v) // 20, 40, 60, 80, 100
    }
}

Синхронный pipeline через Generator:

<?php

declare(strict_types=1);

namespace App\Concurrency;

final readonly class NumberPipeline
{
    /** @param iterable<int> $source */
    public function filter(iterable $source, \Closure $predicate): \Generator
    {
        foreach ($source as $n) {
            if ($predicate($n)) {
                yield $n;
            }
        }
    }

    /** @param iterable<int> $source */
    public function multiply(iterable $source, int $factor): \Generator
    {
        foreach ($source as $n) {
            yield $n * $factor;
        }
    }

    public function run(): \Generator
    {
        $nums = [1, 2, 3, 4, 5, 6, 7, 8, 9, 10];
        $evens = $this->filter($nums, static fn (int $n): bool => $n % 2 === 0);

        return $this->multiply($evens, 10);
    }
}

В PHP pipeline с горутинами в каждой стадии не нужен - in-process данные чаще обрабатываются через ленивые итераторы (Generator) синхронно. Если стадии тяжёлые и нужна параллелизация, каждая стадия становится отдельной очередью в Symfony Messenger: handler одной очереди диспатчит следующее сообщение. Альтернатива - middleware-цепочка внутри одной очереди.

Параллельный pipeline через Messenger - каждая стадия отдельная очередь:

framework:
    messenger:
        transports:
            stage_filter: 'redis://redis:6379/filter'
            stage_multiply: 'redis://redis:6379/multiply'
        routing:
            'App\Pipeline\Stage1Message': stage_filter
            'App\Pipeline\Stage2Message': stage_multiply

Fan-Out / Fan-In

Fan-out: распределяем работу на несколько горутин. Fan-in: собираем результаты в один канал.

Fan-out распыляет один канал на N воркеров, fan-in сливает результаты обратно в один

func fanOut(in <-chan int, workers int) []<-chan int {
    channels := make([]<-chan int, workers)
    for i := 0; i < workers; i++ {
        channels[i] = heavyWork(in)
    }
    return channels
}

func heavyWork(in <-chan int) <-chan int {
    out := make(chan int)
    go func() {
        for n := range in {
            time.Sleep(100 * time.Millisecond) // тяжёлая работа
            out <- n * n
        }
        close(out)
    }()
    return out
}

func fanIn(channels ...<-chan int) <-chan int {
    out := make(chan int)
    var wg sync.WaitGroup

    for _, ch := range channels {
        wg.Add(1)
        go func(c <-chan int) {
            defer wg.Done()
            for v := range c {
                out <- v
            }
        }(ch)
    }

    go func() {
        wg.Wait()
        close(out)
    }()

    return out
}
<?php

declare(strict_types=1);

namespace App\Concurrency;

use Symfony\Component\Messenger\Attribute\AsMessageHandler;
use Symfony\Component\Messenger\MessageBusInterface;

final readonly class WorkItem
{
    public function __construct(public int $value, public string $batchId) {}
}

final readonly class WorkResult
{
    public function __construct(public int $result, public string $batchId) {}
}

#[AsMessageHandler]
final readonly class FanOutHandler
{
    public function __construct(private MessageBusInterface $bus) {}

    public function __invoke(WorkItem $item): void
    {
        \usleep(100_000); // тяжёлая работа
        $result = $item->value * $item->value;

        // fan-in: пишем результат в общий transport
        $this->bus->dispatch(new WorkResult($result, $item->batchId));
    }
}

В PHP fan-out - это N воркеров (отдельных процессов), читающих одну очередь. Брокер сам распределяет сообщения. Fan-in - все воркеры пишут в одну результирующую очередь / БД-таблицу, агрегатор читает оттуда. Координация "когда все закончили" - через счётчик в Redis или БД, а не через WaitGroup.

Конфигурация: одна очередь work_items, N воркеров через supervisor (numprocs=10). Брокер раздаёт сообщения round-robin между ними.

Worker Pool

Worker pool: фиксированное число воркеров обрабатывает jobs канал, ctx.Done для остановки

func workerPool[T any, R any](
    ctx context.Context,
    workers int,
    jobs <-chan T,
    process func(T) R,
) <-chan R {
    results := make(chan R)
    var wg sync.WaitGroup

    for i := 0; i < workers; i++ {
        wg.Add(1)
        go func() {
            defer wg.Done()
            for {
                select {
                case job, ok := <-jobs:
                    if !ok {
                        return
                    }
                    results <- process(job)
                case <-ctx.Done():
                    return
                }
            }
        }()
    }

    go func() {
        wg.Wait()
        close(results)
    }()

    return results
}

Semaphore (ограничение конкурентности)

type Semaphore struct {
    ch chan struct{}
}

func NewSemaphore(max int) *Semaphore {
    return &Semaphore{ch: make(chan struct{}, max)}
}

func (s *Semaphore) Acquire() { s.ch <- struct{}{} }
func (s *Semaphore) Release() { <-s.ch }

// Использование
sem := NewSemaphore(10) // максимум 10 одновременных запросов

for _, url := range urls {
    sem.Acquire()
    go func(u string) {
        defer sem.Release()
        fetch(u)
    }(url)
}
<?php

declare(strict_types=1);

namespace App\Concurrency;

use Symfony\Contracts\HttpClient\HttpClientInterface;

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

    /** @param list<string> $urls */
    public function fetchAll(array $urls): array
    {
        // HttpClient умеет параллелить запросы через stream()
        // max_host_connections в config/packages/framework.yaml ограничит пул
        $responses = [];
        foreach ($urls as $url) {
            $responses[] = $this->client->request('GET', $url);
        }

        $results = [];
        foreach ($this->client->stream($responses) as $response => $chunk) {
            if ($chunk->isLast()) {
                $results[] = $response->getContent();
            }
        }

        return $results;
    }
}

В PHP «семафор» - это либо конкурентность HTTP-клиента (Symfony HttpClient с max_host_connections), либо symfony/lock с counting semaphore-паттерном через Redis. Для длинноживущих воркеров - конфигурация numprocs в supervisor. Реальное N выбирается под backend (БД pool, rate limit API).

Конфигурация лимита параллелизма:

framework:
    http_client:
        default_options:
            max_host_connections: 10 # эквивалент семафора

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