HardПрактика11 min

Паттерны конкурентности

Pipeline, fan-out/fan-in, worker pool, semaphore, rate limiter, errgroup и другие паттерны

Обзор паттернов

Go предоставляет мощные примитивы (горутины, каналы, select), из которых строятся стандартные паттерны конкурентного программирования. Эти паттерны решают типичные задачи: распределение работы, сбор результатов, ограничение параллелизма, обработка ошибок.

┌──────────────────────────────────────────────────────────┐
│                  Паттерны конкурентности                  │
├──────────────┬───────────────┬────────────────────────────┤
│  Потоковые   │  Управление   │  Утилитарные              │
│              │  ресурсами     │                           │
├──────────────┼───────────────┼────────────────────────────┤
│  Pipeline    │  Worker Pool  │  Generator                │
│  Fan-out     │  Semaphore    │  Or-done                  │
│  Fan-in      │  Rate Limiter │  Tee                      │
│  Bridge      │  errgroup     │  singleflight             │
└──────────────┴───────────────┴────────────────────────────┘

Pipeline (Конвейер)

Pipeline -- цепочка стадий обработки, где каждая стадия -- функция, получающая данные из входного канала и отправляющая результат в выходной.

Input → [Stage 1] → [Stage 2] → [Stage 3] → Output
         генерация    обработка    фильтрация
package main

import "fmt"

// generator creates a channel and sends values into it
func generator(nums ...int) <-chan int {
    out := make(chan int)
    go func() {
        defer close(out)
        for _, n := range nums {
            out <- n
        }
    }()
    return out
}

// square reads from in, squares each value, sends to out
func square(in <-chan int) <-chan int {
    out := make(chan int)
    go func() {
        defer close(out)
        for n := range in {
            out <- n * n
        }
    }()
    return out
}

// filter keeps only values matching predicate
func filter(in <-chan int, pred func(int) bool) <-chan int {
    out := make(chan int)
    go func() {
        defer close(out)
        for n := range in {
            if pred(n) {
                out <- n
            }
        }
    }()
    return out
}

func main() {
    // Build pipeline: generate → square → filter (>10)
    nums := generator(1, 2, 3, 4, 5, 6, 7, 8, 9, 10)
    squared := square(nums)
    big := filter(squared, func(n int) bool { return n > 10 })

    // Consume results
    for val := range big {
        fmt.Println(val) // 16, 25, 36, 49, 64, 81, 100
    }
}

Pipeline с отменой через context

package main

import (
    "context"
    "fmt"
)

// generatorCtx respects context cancellation
func generatorCtx(ctx context.Context, nums ...int) <-chan int {
    out := make(chan int)
    go func() {
        defer close(out)
        for _, n := range nums {
            select {
            case <-ctx.Done():
                return
            case out <- n:
            }
        }
    }()
    return out
}

// squareCtx squares values with cancellation support
func squareCtx(ctx context.Context, in <-chan int) <-chan int {
    out := make(chan int)
    go func() {
        defer close(out)
        for n := range in {
            select {
            case <-ctx.Done():
                return
            case out <- n * n:
            }
        }
    }()
    return out
}

func main() {
    ctx, cancel := context.WithCancel(context.Background())
    defer cancel()

    nums := generatorCtx(ctx, 1, 2, 3, 4, 5, 6, 7, 8, 9, 10)
    squared := squareCtx(ctx, nums)

    // Take only first 3 results, then cancel pipeline
    count := 0
    for val := range squared {
        fmt.Println(val)
        count++
        if count == 3 {
            cancel() // stops all stages
            break
        }
    }
}

Fan-out, Fan-in

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

                    ┌── Worker 1 ──┐
                    │              │
Input ──► Fan-out ──┼── Worker 2 ──┼── Fan-in ──► Output
                    │              │
                    └── Worker 3 ──┘
package main

import (
    "fmt"
    "sync"
    "time"
)

// producer generates work items
func producer(count int) <-chan int {
    out := make(chan int)
    go func() {
        defer close(out)
        for i := 1; i <= count; i++ {
            out <- i
        }
    }()
    return out
}

// worker processes items (fan-out — multiple workers read from same channel)
func worker(id int, in <-chan int) <-chan string {
    out := make(chan string)
    go func() {
        defer close(out)
        for n := range in {
            time.Sleep(50 * time.Millisecond) // simulate work
            out <- fmt.Sprintf("worker %d processed %d", id, n)
        }
    }()
    return out
}

// merge combines multiple channels into one (fan-in)
func merge(channels ...<-chan string) <-chan string {
    var wg sync.WaitGroup
    merged := make(chan string)

    // Start a goroutine for each input channel
    for _, ch := range channels {
        wg.Add(1)
        go func() {
            defer wg.Done()
            for val := range ch {
                merged <- val
            }
        }()
    }

    // Close merged channel when all input channels are done
    go func() {
        wg.Wait()
        close(merged)
    }()

    return merged
}

func main() {
    // Generate work
    jobs := producer(20)

    // Fan-out: 4 workers reading from same channel
    w1 := worker(1, jobs)
    w2 := worker(2, jobs)
    w3 := worker(3, jobs)
    w4 := worker(4, jobs)

    // Fan-in: merge all results
    results := merge(w1, w2, w3, w4)

    // Consume
    for result := range results {
        fmt.Println(result)
    }
}

Worker Pool (Пул воркеров)

Worker Pool -- фиксированное число горутин, обрабатывающих задачи из общего канала. Это самый распространённый паттерн для ограничения параллелизма.

                     ┌─────────┐
              ┌──────│ Worker 1│──────┐
              │      └─────────┘      │
┌──────────┐  │      ┌─────────┐      │  ┌───────────┐
│  Jobs    ├──┼──────│ Worker 2│──────┼──│  Results  │
│ (chan)   │  │      └─────────┘      │  │  (chan)    │
└──────────┘  │      ┌─────────┐      │  └───────────┘
              └──────│ Worker 3│──────┘
                     └─────────┘
package main

import (
    "fmt"
    "math/rand"
    "sync"
    "time"
)

// Job represents a unit of work
type Job struct {
    ID      int
    Payload string
}

// Result represents the output of processing a job
type Result struct {
    Job      Job
    Output   string
    Duration time.Duration
}

// workerPool processes jobs from the jobs channel
func workerPool(id int, jobs <-chan Job, results chan<- Result, wg *sync.WaitGroup) {
    defer wg.Done()
    for job := range jobs {
        start := time.Now()

        // Simulate variable processing time
        time.Sleep(time.Duration(rand.Intn(100)) * time.Millisecond)

        results <- Result{
            Job:      job,
            Output:   fmt.Sprintf("worker-%d processed: %s", id, job.Payload),
            Duration: time.Since(start),
        }
    }
}

func main() {
    const (
        numWorkers = 5
        numJobs    = 20
    )

    jobs := make(chan Job, numJobs)
    results := make(chan Result, numJobs)

    // Start workers
    var wg sync.WaitGroup
    for w := 1; w <= numWorkers; w++ {
        wg.Add(1)
        go workerPool(w, jobs, results, &wg)
    }

    // Send jobs
    for j := 1; j <= numJobs; j++ {
        jobs <- Job{ID: j, Payload: fmt.Sprintf("task-%d", j)}
    }
    close(jobs) // signal no more jobs

    // Close results when all workers done
    go func() {
        wg.Wait()
        close(results)
    }()

    // Collect results
    for result := range results {
        fmt.Printf("[%v] %s\n", result.Duration.Round(time.Millisecond), result.Output)
    }
}

Semaphore (Семафор)

Семафор ограничивает количество одновременно выполняющихся операций. Реализуется через буферизованный канал.

package main

import (
    "fmt"
    "sync"
    "time"
)

// Semaphore limits concurrent access
type Semaphore chan struct{}

// NewSemaphore creates a semaphore with max concurrent count
func NewSemaphore(max int) Semaphore {
    return make(Semaphore, max)
}

// Acquire takes a slot (blocks if full)
func (s Semaphore) Acquire() {
    s <- struct{}{}
}

// Release frees a slot
func (s Semaphore) Release() {
    <-s
}

func main() {
    // Allow max 3 concurrent operations
    sem := NewSemaphore(3)
    var wg sync.WaitGroup

    for i := 1; i <= 10; i++ {
        wg.Add(1)
        go func() {
            defer wg.Done()

            sem.Acquire()
            defer sem.Release()

            fmt.Printf("Task %d: started (concurrent slots used: %d/%d)\n",
                i, len(sem), cap(sem))
            time.Sleep(500 * time.Millisecond) // simulate work
            fmt.Printf("Task %d: done\n", i)
        }()
    }

    wg.Wait()
    fmt.Println("All tasks completed")
}

Generator (Генератор)

Генератор -- функция, возвращающая канал только для чтения. Горутина генерирует значения и отправляет в канал.

package main

import "fmt"

// fibonacci returns a channel that produces Fibonacci numbers
func fibonacci(n int) <-chan int {
    ch := make(chan int)
    go func() {
        defer close(ch)
        a, b := 0, 1
        for range n {
            ch <- a
            a, b = b, a+b
        }
    }()
    return ch
}

// repeat generates value infinitely until context cancelled
func repeat(done <-chan struct{}, values ...int) <-chan int {
    out := make(chan int)
    go func() {
        defer close(out)
        for {
            for _, v := range values {
                select {
                case <-done:
                    return
                case out <- v:
                }
            }
        }
    }()
    return out
}

// take reads n values from a channel
func take(done <-chan struct{}, in <-chan int, n int) <-chan int {
    out := make(chan int)
    go func() {
        defer close(out)
        for range n {
            select {
            case <-done:
                return
            case val, ok := <-in:
                if !ok {
                    return
                }
                out <- val
            }
        }
    }()
    return out
}

func main() {
    // Fibonacci generator
    fmt.Println("Fibonacci:")
    for n := range fibonacci(10) {
        fmt.Printf("%d ", n)
    }
    fmt.Println()

    // Infinite repeat with take
    done := make(chan struct{})
    stream := repeat(done, 1, 2, 3)
    limited := take(done, stream, 7)

    fmt.Println("Repeat-take:")
    for n := range limited {
        fmt.Printf("%d ", n) // 1 2 3 1 2 3 1
    }
    fmt.Println()
    close(done)
}

Or-done канал

Паттерн or-done оборачивает канал, добавляя поддержку отмены через done-канал:

package main

import "fmt"

// orDone wraps a channel to support cancellation via done
func orDone(done <-chan struct{}, in <-chan int) <-chan int {
    out := make(chan int)
    go func() {
        defer close(out)
        for {
            select {
            case <-done:
                return
            case val, ok := <-in:
                if !ok {
                    return
                }
                select {
                case <-done:
                    return
                case out <- val:
                }
            }
        }
    }()
    return out
}

func main() {
    done := make(chan struct{})
    source := make(chan int, 10)

    // Fill source
    for i := range 10 {
        source <- i
    }
    close(source)

    // Read with or-done (graceful cancellation)
    for val := range orDone(done, source) {
        fmt.Println(val)
        if val == 5 {
            close(done) // cancel after 5
            break
        }
    }
}

Tee канал

Tee разделяет один канал на два, отправляя каждое значение в оба выходных канала:

package main

import "fmt"

// tee splits one channel into two
func tee(done <-chan struct{}, in <-chan int) (<-chan int, <-chan int) {
    out1 := make(chan int)
    out2 := make(chan int)

    go func() {
        defer close(out1)
        defer close(out2)

        for val := range orDoneInline(done, in) {
            // Use local copies for select to avoid double-send
            ch1, ch2 := out1, out2
            for range 2 {
                select {
                case ch1 <- val:
                    ch1 = nil // disable after send
                case ch2 <- val:
                    ch2 = nil // disable after send
                }
            }
        }
    }()

    return out1, out2
}

// orDoneInline is an inline version of or-done
func orDoneInline(done <-chan struct{}, in <-chan int) <-chan int {
    out := make(chan int)
    go func() {
        defer close(out)
        for {
            select {
            case <-done:
                return
            case v, ok := <-in:
                if !ok {
                    return
                }
                select {
                case <-done:
                    return
                case out <- v:
                }
            }
        }
    }()
    return out
}

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

    source := make(chan int, 5)
    for i := 1; i <= 5; i++ {
        source <- i
    }
    close(source)

    ch1, ch2 := tee(done, source)

    // Read from both channels
    for val := range ch1 {
        fmt.Printf("ch1: %d, ch2: %d\n", val, <-ch2)
    }
}

Bridge канал

Bridge объединяет канал каналов в один плоский канал:

package main

import "fmt"

// bridge flattens a channel of channels into a single channel
func bridge(done <-chan struct{}, chanStream <-chan <-chan int) <-chan int {
    out := make(chan int)
    go func() {
        defer close(out)
        for {
            var stream <-chan int
            select {
            case <-done:
                return
            case maybeStream, ok := <-chanStream:
                if !ok {
                    return
                }
                stream = maybeStream
            }

            for val := range stream {
                select {
                case <-done:
                    return
                case out <- val:
                }
            }
        }
    }()
    return out
}

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

    // Create channel of channels
    genBatch := func(begin, end int) <-chan int {
        ch := make(chan int)
        go func() {
            defer close(ch)
            for i := begin; i <= end; i++ {
                ch <- i
            }
        }()
        return ch
    }

    chanStream := make(chan (<-chan int), 3)
    go func() {
        defer close(chanStream)
        chanStream <- genBatch(1, 3)
        chanStream <- genBatch(4, 6)
        chanStream <- genBatch(7, 9)
    }()

    // Bridge flattens: [1,2,3], [4,5,6], [7,8,9] → 1,2,3,4,5,6,7,8,9
    for val := range bridge(done, chanStream) {
        fmt.Printf("%d ", val)
    }
    fmt.Println()
}

Rate Limiter (Ограничитель скорости)

package main

import (
    "fmt"
    "time"
)

func main() {
    // Basic rate limiter: 1 request per 200ms
    limiter := time.NewTicker(200 * time.Millisecond)
    defer limiter.Stop()

    requests := []int{1, 2, 3, 4, 5}

    for _, req := range requests {
        <-limiter.C // wait for tick
        fmt.Printf("[%s] Processing request %d\n",
            time.Now().Format("15:04:05.000"), req)
    }

    fmt.Println("\n--- Burst rate limiter ---")

    // Burst rate limiter: allow burst of 3, then 1 per 200ms
    burstyLimiter := make(chan time.Time, 3)

    // Fill initial burst
    for range 3 {
        burstyLimiter <- time.Now()
    }

    // Refill at steady rate
    go func() {
        for t := range time.Tick(200 * time.Millisecond) {
            burstyLimiter <- t
        }
    }()

    burstyRequests := []int{1, 2, 3, 4, 5, 6, 7}
    for _, req := range burstyRequests {
        <-burstyLimiter
        fmt.Printf("[%s] Bursty request %d\n",
            time.Now().Format("15:04:05.000"), req)
    }
}

Обработка ошибок: errgroup

Пакет golang.org/x/sync/errgroup -- стандартный инструмент для запуска группы горутин с обработкой ошибок:

package main

import (
    "context"
    "fmt"
    "math/rand"
    "time"

    "golang.org/x/sync/errgroup"
)

// fetchURL simulates fetching a URL
func fetchURL(ctx context.Context, url string) (string, error) {
    // Simulate network delay
    delay := time.Duration(rand.Intn(500)) * time.Millisecond

    select {
    case <-ctx.Done():
        return "", ctx.Err()
    case <-time.After(delay):
        if rand.Float32() < 0.1 {
            return "", fmt.Errorf("fetch %s: connection timeout", url)
        }
        return fmt.Sprintf("Content of %s", url), nil
    }
}

func main() {
    urls := []string{
        "https://go.dev",
        "https://pkg.go.dev",
        "https://blog.golang.org",
        "https://tour.golang.org",
    }

    // errgroup with context — if any goroutine fails, context is cancelled
    g, ctx := errgroup.WithContext(context.Background())

    results := make([]string, len(urls))

    for i, url := range urls {
        g.Go(func() error {
            content, err := fetchURL(ctx, url)
            if err != nil {
                return err
            }
            results[i] = content // safe: each goroutine writes to unique index
            return nil
        })
    }

    // Wait for all goroutines and check for errors
    if err := g.Wait(); err != nil {
        fmt.Println("Error:", err)
        return
    }

    for _, r := range results {
        fmt.Println(r)
    }
}

errgroup с ограничением параллелизма

package main

import (
    "fmt"
    "time"

    "golang.org/x/sync/errgroup"
)

func main() {
    g := new(errgroup.Group)

    // Limit to 3 concurrent goroutines
    g.SetLimit(3)

    for i := 1; i <= 10; i++ {
        g.Go(func() error {
            fmt.Printf("Task %d started\n", i)
            time.Sleep(500 * time.Millisecond)
            fmt.Printf("Task %d done\n", i)
            return nil
        })
    }

    if err := g.Wait(); err != nil {
        fmt.Println("Error:", err)
    }
}

golang.org/x/sync/semaphore

Взвешенный семафор из стандартной расширенной библиотеки:

package main

import (
    "context"
    "fmt"
    "sync"
    "time"

    "golang.org/x/sync/semaphore"
)

func main() {
    // Allow max 3 concurrent heavy operations
    sem := semaphore.NewWeighted(3)
    ctx := context.Background()

    var wg sync.WaitGroup
    for i := 1; i <= 10; i++ {
        wg.Add(1)
        go func() {
            defer wg.Done()

            // Acquire 1 weight (blocks if 3 are already acquired)
            if err := sem.Acquire(ctx, 1); err != nil {
                fmt.Printf("Task %d: acquire error: %v\n", i, err)
                return
            }
            defer sem.Release(1)

            fmt.Printf("Task %d: running\n", i)
            time.Sleep(300 * time.Millisecond)
            fmt.Printf("Task %d: done\n", i)
        }()
    }

    wg.Wait()
}

golang.org/x/sync/singleflight

singleflight подавляет дублирующиеся запросы. Если несколько горутин запрашивают одни и те же данные, выполняется только один запрос, а результат разделяется.

package main

import (
    "fmt"
    "sync"
    "time"

    "golang.org/x/sync/singleflight"
)

var sf singleflight.Group

// expensiveLookup simulates a costly operation (DB query, API call)
func expensiveLookup(key string) (string, error) {
    fmt.Printf("  [DB] Actually querying for key: %s\n", key)
    time.Sleep(200 * time.Millisecond)
    return fmt.Sprintf("value_for_%s", key), nil
}

// cachedLookup deduplicates concurrent requests for the same key
func cachedLookup(key string) (string, error) {
    result, err, shared := sf.Do(key, func() (any, error) {
        return expensiveLookup(key)
    })
    if err != nil {
        return "", err
    }
    if shared {
        fmt.Printf("  [Cache] Result was shared for key: %s\n", key)
    }
    return result.(string), nil
}

func main() {
    var wg sync.WaitGroup

    // 10 goroutines requesting the SAME key simultaneously
    for i := range 10 {
        wg.Add(1)
        go func() {
            defer wg.Done()
            val, err := cachedLookup("user:123")
            if err != nil {
                fmt.Printf("Goroutine %d: error: %v\n", i, err)
                return
            }
            fmt.Printf("Goroutine %d: got %s\n", i, val)
        }()
    }

    wg.Wait()
    // "[DB] Actually querying..." prints ONCE, result shared to all 10
}

Применение singleflight: Идеален для кешей, DNS-резолверов, дедупликации API-запросов. Предотвращает "thundering herd" -- лавину одинаковых запросов при протухании кеша.

Проверь себя

Что такое fan-out в контексте конкурентного программирования Go?

Какой пакет из golang.org/x/sync предотвращает дублирующиеся запросы?

Как ограничить количество одновременно работающих горутин в errgroup?

Для чего паттерн bridge-channel?

Code Challenges

Пул воркеров

GO

Реализуйте паттерн Worker Pool с использованием каналов. Функция принимает задачи и количество воркеров, возвращает результаты.

Test Cases

1. Input: tasks=[1,2,3,4,5], workers=2, fn=double→ Expected: [2,4,6,8,10]
2. Input: tasks=[10,20], workers=3, fn=square→ Expected: [100,400]
3. Input: tasks=[], workers=2, fn=identity→ Expected: []