Обзор паттернов
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" -- лавину одинаковых запросов при протухании кеша.