Паттерны: fan-in, fan-out, pipeline, worker pool
Паттерны: fan-in, fan-out, pipeline, worker pool
Эти паттерны - строительные блоки для конкурентных систем поверх горутин и каналов.
Pipeline
Последовательная обработка данных через цепочку этапов.
// Генератор → Фильтр → Маппер → Потребитель
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: собираем результаты в один канал.
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
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 # эквивалент семафора