Go Worker Pool и конкурентная обработка задач: паттерны для высоконагруженных систем
Введение: зачем нужен worker pool в Go-сервисах
В 2026 году Go остаётся одним из доминирующих языков для построения высоконагруженных backend-систем. Горутины дёшевы — их можно запустить миллионы, — но это не означает, что бесконтрольный spawn горутин является правильной стратегией. В реальных сервисах неограниченный параллелизм приводит к исчерпанию файловых дескрипторов, давлению на сборщик мусора, перегрузке downstream-зависимостей (баз данных, сторонних API) и непредсказуемым задержкам.
Worker pool решает ключевые задачи:
- Ограничение параллелизма — не больше N задач выполняется одновременно.
- Переиспользование горутин — вместо создания горутины на каждый запрос используется фиксированный пул.
- Backpressure — при переполнении очереди производитель блокируется или получает ошибку, вместо того чтобы бесконтрольно генерировать нагрузку.
- Изоляция ошибок — паника в одном воркере не роняет всю программу при наличии корректного recover.
В этой статье мы последовательно разберём всё: от примитивной реализации до production-ready решения с динамическим масштабированием, retry-логикой, graceful shutdown, интеграцией с PostgreSQL в качестве надёжной очереди и полноценным мониторингом.
Базовая реализация worker pool на горутинах и каналах
Классический worker pool в Go строится на трёх примитивах: канал задач (jobs chan), канал результатов (results chan) и sync.WaitGroup для ожидания завершения всех воркеров.
package workerpool
import (
"context"
"fmt"
"sync"
)
// Job представляет единицу работы, которую нужно выполнить.
type Job struct {
ID int
Payload any
}
// Result хранит результат выполнения задачи.
type Result struct {
JobID int
Value any
Err error
}
// ProcessFunc — функция обработки задачи, передаётся при создании пула.
type ProcessFunc func(ctx context.Context, job Job) Result
// Pool — базовый worker pool.
type Pool struct {
jobs chan Job
results chan Result
wg sync.WaitGroup
process ProcessFunc
size int
}
// NewPool создаёт пул с size воркерами и буфером очереди queueSize.
func NewPool(size, queueSize int, fn ProcessFunc) *Pool {
return &Pool{
jobs: make(chan Job, queueSize),
results: make(chan Result, queueSize),
process: fn,
size: size,
}
}
// Start запускает воркеры и начинает обработку задач.
func (p *Pool) Start(ctx context.Context) {
for i := 0; i < p.size; i++ {
p.wg.Add(1)
go p.worker(ctx)
}
}
// worker — горутина, читающая задачи из канала jobs.
func (p *Pool) worker(ctx context.Context) {
defer p.wg.Done()
for {
select {
case job, ok := <-p.jobs:
if !ok {
// Канал закрыт — воркер завершает работу.
return
}
result := p.process(ctx, job)
// Отправляем результат; если канал переполнен — блокируемся.
select {
case p.results <- result:
case <-ctx.Done():
return
}
case <-ctx.Done():
return
}
}
}
// Submit отправляет задачу в очередь.
// Возвращает ошибку, если контекст отменён или канал переполнен.
func (p *Pool) Submit(ctx context.Context, job Job) error {
select {
case p.jobs <- job:
return nil
case <-ctx.Done():
return fmt.Errorf("worker pool: submit cancelled: %w", ctx.Err())
}
}
// Results возвращает канал результатов для чтения потребителем.
func (p *Pool) Results() <-chan Result {
return p.results
}
// Stop закрывает канал задач и ждёт завершения всех воркеров.
func (p *Pool) Stop() {
close(p.jobs)
p.wg.Wait()
close(p.results)
}
Ключевые архитектурные решения в этой реализации:
- Канал
jobsбуферизован — производитель не блокируется сразу при кратковременных всплесках нагрузки. - Воркер слушает сразу два канала через
select: задачи и отмену контекста — это предотвращает зависание при shutdown. - Закрытие
jobsявляется сигналом воркерам о завершении работы: Go гарантирует, что после закрытия канала все ранее отправленные значения будут прочитаны.
Динамическое масштабирование пула: адаптация под нагрузку
Статический пул оптимален при предсказуемой нагрузке. В реальных системах нагрузка меняется: ночью трафик минимален, в часы пик — в 10 раз выше. Динамический пул позволяет масштабировать количество воркеров в диапазоне [minWorkers, maxWorkers] на основе метрики — длины очереди задач.
package workerpool
import (
"context"
"sync"
"sync/atomic"
"time"
)
// DynamicPool масштабирует количество воркеров в зависимости от нагрузки.
type DynamicPool struct {
jobs chan Job
results chan Result
process ProcessFunc
mu sync.Mutex
wg sync.WaitGroup
ctx context.Context
cancel context.CancelFunc
active atomic.Int64 // текущее количество воркеров
minWorkers int
maxWorkers int
}
func NewDynamicPool(min, max, queueSize int, fn ProcessFunc) *DynamicPool {
ctx, cancel := context.WithCancel(context.Background())
p := &DynamicPool{
jobs: make(chan Job, queueSize),
results: make(chan Result, queueSize),
process: fn,
ctx: ctx,
cancel: cancel,
minWorkers: min,
maxWorkers: max,
}
// Запускаем минимальное количество воркеров сразу.
for i := 0; i < min; i++ {
p.startWorker()
}
// Запускаем горутину-автоскейлер.
go p.autoscale()
return p
}
func (p *DynamicPool) startWorker() {
p.active.Add(1)
p.wg.Add(1)
go func() {
defer func() {
p.active.Add(-1)
p.wg.Done()
}()
for {
select {
case job, ok := <-p.jobs:
if !ok {
return
}
result := p.process(p.ctx, job)
select {
case p.results <- result:
case <-p.ctx.Done():
return
}
case <-p.ctx.Done():
return
}
}
}()
}
// autoscale проверяет нагрузку каждые 500 мс и при необходимости добавляет воркеров.
func (p *DynamicPool) autoscale() {
ticker := time.NewTicker(500 * time.Millisecond)
defer ticker.Stop()
for {
select {
case <-ticker.C:
queueLen := len(p.jobs)
current := int(p.active.Load())
// Если очередь заполнена более чем на 50% и можно расти — добавляем воркера.
if queueLen > cap(p.jobs)/2 && current < p.maxWorkers {
p.startWorker()
}
// Если очередь пуста и воркеров больше минимума — останавливаем одного.
// Для остановки отдельного воркера используем drain-канал (упрощённо).
case <-p.ctx.Done():
return
}
}
}
func (p *DynamicPool) Stop() {
p.cancel()
close(p.jobs)
p.wg.Wait()
close(p.results)
}
Для уменьшения числа воркеров в production-реализациях используют отдельный quit-канал на каждого воркера или библиотеку golang.org/x/sync/semaphore для точного контроля над параллелизмом.
Обработка ошибок и retry-логика без утечек горутин
В высоконагруженных системах часть задач неизбежно завершается ошибкой: сетевые сбои, временная недоступность зависимостей, дедлайны контекста. Наивный retry в цикле внутри воркера блокирует его и снижает пропускную способность пула. Правильный подход — повторная отправка задачи в очередь с экспоненциальной задержкой и ограниченным числом попыток.
package workerpool
import (
"context"
"errors"
"log/slog"
"math"
"time"
)
// RetryableJob расширяет Job метаданными для retry.
type RetryableJob struct {
Job
Attempt int
MaxRetries int
NextRunAt time.Time
}
// withRetry оборачивает ProcessFunc, добавляя retry-логику.
func withRetry(pool *Pool, fn ProcessFunc) ProcessFunc {
return func(ctx context.Context, job Job) Result {
rj, ok := job.Payload.(RetryableJob)
if !ok {
return fn(ctx, job)
}
result := fn(ctx, job)
if result.Err == nil {
return result
}
// Не повторяем при отмене контекста.
if errors.Is(result.Err, context.Canceled) || errors.Is(result.Err, context.DeadlineExceeded) {
return result
}
if rj.Attempt >= rj.MaxRetries {
slog.Error("job failed permanently",
"job_id", job.ID,
"attempts", rj.Attempt,
"err", result.Err,
)
return result
}
// Экспоненциальная задержка: 2^attempt секунд, максимум 60 секунд.
delay := time.Duration(math.Min(math.Pow(2, float64(rj.Attempt)), 60)) * time.Second
rj.Attempt++
rj.NextRunAt = time.Now().Add(delay)
// Планируем повторную отправку в отдельной горутине.
// Горутина завершится через delay или при отмене контекста — утечки нет.
go func() {
select {
case <-time.After(delay):
retryJob := Job{ID: job.ID, Payload: rj}
if err := pool.Submit(ctx, retryJob); err != nil {
slog.Warn("failed to resubmit job", "job_id", job.ID, "err", err)
}
case <-ctx.Done():
slog.Info("retry cancelled due to context", "job_id", job.ID)
}
}()
return Result{JobID: job.ID, Err: nil} // не считаем временную ошибку финальной
}
}
Критически важный момент: горутина для отложенного retry всегда завершается — либо по таймеру, либо при отмене контекста. Это исключает утечку горутин (goroutine leak) — одну из самых опасных проблем в долгоживущих Go-сервисах.
Graceful Shutdown: корректное завершение всех задач
Production-сервис должен корректно завершать работу при получении SIGTERM или SIGINT: дообработать уже принятые задачи, не принимать новые и освободить ресурсы.
package main
import (
"context"
"log/slog"
"os"
"os/signal"
"syscall"
"time"
"github.com/yourorg/workerpool"
)
func main() {
// Контекст приложения: отменяется при получении OS-сигнала.
ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGTERM, syscall.SIGINT)
defer stop()
pool := workerpool.NewPool(10, 100, processImage)
pool.Start(ctx)
// Горутина-продюсер: останавливается при отмене контекста.
go func() {
for i := 0; ; i++ {
select {
case <-ctx.Done():
slog.Info("producer stopped")
return
default:
}
job := workerpool.Job{ID: i, Payload: fmt.Sprintf("image_%d.jpg", i)}
if err := pool.Submit(ctx, job); err != nil {
slog.Warn("submit failed", "err", err)
}
}
}()
// Читаем результаты в отдельной горутине.
go func() {
for result := range pool.Results() {
if result.Err != nil {
slog.Error("job error", "job_id", result.JobID, "err", result.Err)
}
}
}()
// Ждём сигнала завершения.
<-ctx.Done()
slog.Info("shutdown signal received, draining pool...")
// Даём воркерам до 30 секунд на дообработку задач.
shutdownCtx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
defer cancel()
done := make(chan struct{})
go func() {
pool.Stop() // закрывает jobs-канал и ждёт WaitGroup
close(done)
}()
select {
case <-done:
slog.Info("graceful shutdown complete")
case <-shutdownCtx.Done():
slog.Error("shutdown timeout exceeded, forcing exit")
os.Exit(1)
}
}
Паттерн signal.NotifyContext появился в Go 1.16 и является идиоматическим способом связать жизненный цикл приложения с OS-сигналами. Таймаут на shutdown (здесь 30 секунд) предотвращает бесконечное ожидание при зависших задачах.
Интеграция с PostgreSQL: надёжная очередь задач
Для задач, которые нельзя потерять при перезапуске сервиса (отправка email, обработка платежей), канальная очередь недостаточна — данные в памяти не персистентны. PostgreSQL с паттерном FOR UPDATE SKIP LOCKED решает эту проблему, превращая таблицу в надёжную распределённую очередь.
-- Схема таблицы задач
CREATE TABLE jobs (
id BIGSERIAL PRIMARY KEY,
type TEXT NOT NULL,
payload JSONB NOT NULL,
status TEXT NOT NULL DEFAULT 'pending', -- pending | running | done | failed
attempts INT NOT NULL DEFAULT 0,
max_retries INT NOT NULL DEFAULT 3,
run_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
);
CREATE INDEX idx_jobs_status_run_at ON jobs (status, run_at)
WHERE status IN ('pending', 'failed');
package pgqueue
import (
"context"
"encoding/json"
"time"
"github.com/jackc/pgx/v5"
"github.com/jackc/pgx/v5/pgxpool"
)
type PGJob struct {
ID int64
Type string
Payload json.RawMessage
Attempts int
MaxRetries int
}
type Queue struct {
db *pgxpool.Pool
}
func NewQueue(db *pgxpool.Pool) *Queue {
return &Queue{db: db}
}
// Dequeue атомарно захватывает одну задачу из очереди.
// FOR UPDATE SKIP LOCKED гарантирует, что конкурентные воркеры не заберут одну задачу дважды.
func (q *Queue) Dequeue(ctx context.Context) (*PGJob, error) {
tx, err := q.db.Begin(ctx)
if err != nil {
return nil, err
}
defer tx.Rollback(ctx)
var job PGJob
err = tx.QueryRow(ctx, `
UPDATE jobs
SET status = 'running',
attempts = attempts + 1,
updated_at = NOW()
WHERE id = (
SELECT id FROM jobs
WHERE status IN ('pending', 'failed')
AND run_at <= NOW()
AND attempts < max_retries
ORDER BY run_at ASC
FOR UPDATE SKIP LOCKED
LIMIT 1
)
RETURNING id, type, payload, attempts, max_retries
`).Scan(&job.ID, &job.Type, &job.Payload, &job.Attempts, &job.MaxRetries)
if err != nil {
if err == pgx.ErrNoRows {
return nil, nil // очередь пуста
}
return nil, err
}
return &job, tx.Commit(ctx)
}
// Complete помечает задачу как успешно выполненную.
func (q *Queue) Complete(ctx context.Context, jobID int64) error {
_, err := q.db.Exec(ctx,
`UPDATE jobs SET status = 'done', updated_at = NOW() WHERE id = $1`,
jobID,
)
return err
}
// Fail помечает задачу как проваленную и планирует retry с экспоненциальной задержкой.
func (q *Queue) Fail(ctx context.Context, jobID int64, attempts int) error {
delay := time.Duration(1<<attempts) * time.Second // 2^attempts секунд
if delay > 10*time.Minute {
delay = 10 * time.Minute
}
_, err := q.db.Exec(ctx, `
UPDATE jobs
SET status = 'failed',
run_at = NOW() + $1::interval,
updated_at = NOW()
WHERE id = $2
`, delay.String(), jobID)
return err
}
Ключевое свойство FOR UPDATE SKIP LOCKED: строки, заблокированные другими транзакциями, пропускаются, а не вызывают блокировку. Это делает паттерн масштабируемым — десятки воркеров могут конкурентно опрашивать одну таблицу без взаимных блокировок.
Для polling-цикла воркера используйте адаптивные паузы: если очередь пуста — увеличивайте интервал опроса (например, от 100 мс до 5 секунд), при появлении задач — возвращайтесь к минимальному интервалу.
Мониторинг worker pool: Prometheus и Grafana
Без метрик worker pool — чёрный ящик. Ключевые показатели для мониторинга: размер очереди, количество активных воркеров, время обработки задачи, частота ошибок и retry.
package metrics
import (
"time"
"github.com/prometheus/client_golang/prometheus"
"github.com/prometheus/client_golang/prometheus/promauto"
)
var (
// Текущая длина очереди задач.
QueueDepth = promauto.NewGauge(prometheus.GaugeOpts{
Name: "worker_pool_queue_depth",
Help: "Number of jobs waiting in the queue.",
})
// Количество активных воркеров.
ActiveWorkers = promauto.NewGauge(prometheus.GaugeOpts{
Name: "worker_pool_active_workers",
Help: "Number of currently active workers.",
})
// Гистограмма времени обработки задачи.
JobDuration = promauto.NewHistogramVec(
prometheus.HistogramOpts{
Name: "worker_pool_job_duration_seconds",
Help: "Time spent processing a job.",
Buckets: prometheus.ExponentialBuckets(0.001, 2, 15), // 1ms..16s
},
[]string{"job_type", "status"}, // status: success | error
)
// Счётчик выполненных задач.
JobsTotal = promauto.NewCounterVec(
prometheus.CounterOpts{
Name: "worker_pool_jobs_total",
Help: "Total number of processed jobs.",
},
[]string{"job_type", "status"},
)
// Счётчик retry.
RetriesTotal = promauto.NewCounter(prometheus.CounterOpts{
Name: "worker_pool_retries_total",
Help: "Total number of job retries.",
})
)
// RecordJobExecution инструментирует выполнение задачи.
func RecordJobExecution(jobType string, start time.Time, err error) {
duration := time.Since(start).Seconds()
status := "success"
if err != nil {
status = "error"
}
JobDuration.WithLabelValues(jobType, status).Observe(duration)
JobsTotal.WithLabelValues(jobType, status).Inc()
}
Рекомендуемые алерты в Grafana:
worker_pool_queue_depth > 1000на протяжении 5 минут — признак перегрузки или зависших воркеров.rate(worker_pool_jobs_total{status="error"}[5m]) / rate(worker_pool_jobs_total[5m]) > 0.05— error rate выше 5%.histogram_quantile(0.99, worker_pool_job_duration_seconds) > 30— p99 latency превышает 30 секунд.
Сравнение с готовыми решениями: когда писать самому
Перед реализацией собственного worker pool стоит оценить готовые решения:
- asynq — Redis-based очередь с rich UI (Asynq Monitor), поддержкой расписаний, приоритетов и уникальных задач. Отличный выбор, если Redis уже есть в стеке и нужна быстрая интеграция.
- river — PostgreSQL-native очередь задач от авторов pgx. Использует
LISTEN/NOTIFYвместо polling, поддерживает транзакционный enqueue (задача добавляется атомарно с бизнес-транзакцией). Идеален, если PostgreSQL — единственная зависимость. - machinery — более тяжёлое решение с поддержкой нескольких брокеров (Redis, AMQP, SQS), хорошо подходит для распределённых workflow.
Когда имеет смысл писать собственный worker pool:
- Задачи in-memory, персистентность не требуется (например, pipeline обработки запросов).
- Нужна специфическая логика приоритезации или маршрутизации задач.
- Минимизация внешних зависимостей критична (embedded-сервисы, edge).
- Готовые решения избыточны по функциональности и добавляют нежелательную сложность.
Практический пример: сервис асинхронной обработки изображений
Соберём всё вместе в реалистичный пример: HTTP-сервис принимает загрузки изображений, помещает задачи обработки (изменение размера, конвертация форматов) в PostgreSQL-очередь, worker pool забирает задачи и обрабатывает их конкурентно.
package main
import (
"context"
"encoding/json"
"log/slog"
"net/http"
"os"
"os/signal"
"syscall"
"time"
"github.com/jackc/pgx/v5/pgxpool"
"github.com/prometheus/client_golang/prometheus/promhttp"
"github.com/yourorg/metrics"
"github.com/yourorg/pgqueue"
)
type ImagePayload struct {
S3Key string `json:"s3_key"`
Widths []int `json:"widths"`
Formats []string `json:"formats"`
}
// processImageJob — бизнес-логика обработки одного изображения.
func processImageJob(ctx context.Context, payload json.RawMessage) error {
var p ImagePayload
if err := json.Unmarshal(payload, &p); err != nil {
return err
}
// Здесь: скачать из S3, изменить размер, сохранить обратно.
// Для примера эмулируем работу.
slog.Info("processing image", "s3_key", p.S3Key, "widths", p.Widths)
time.Sleep(100 * time.Millisecond) // имитация работы
return nil
}
func runWorkers(ctx context.Context, queue *pgqueue.Queue, workerCount int) {
sem := make(chan struct{}, workerCount) // семафор на число воркеров
for {
select {
case <-ctx.Done():
// Ждём освобождения всех слотов семафора.
for i := 0; i < workerCount; i++ {
sem <- struct{}{}
}
return
default:
}
job, err := queue.Dequeue(ctx)
if err != nil {
slog.Error("dequeue error", "err", err)
time.Sleep(time.Second)
continue
}
if job == nil {
// Очередь пуста — адаптивная пауза.
time.Sleep(200 * time.Millisecond)
continue
}
// Занимаем слот семафора.
sem <- struct{}{}
metrics.ActiveWorkers.Inc()
go func(j *pgqueue.PGJob) {
defer func() {
<-sem // освобождаем слот
metrics.ActiveWorkers.Dec()
}()
start := time.Now()
err := processImageJob(ctx, j.Payload)
metrics.RecordJobExecution(j.Type, start, err)
if err != nil {
slog.Error("job failed", "job_id", j.ID, "err", err)
if qErr := queue.Fail(ctx, j.ID, j.Attempts); qErr != nil {
slog.Error("failed to mark job as failed", "err", qErr)
}
return
}
if err := queue.Complete(ctx, j.ID); err != nil {
slog.Error("failed to mark job complete", "err", err)
}
}(job)
}
}
func main() {
ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGTERM, syscall.SIGINT)
defer stop()
db, err := pgxpool.New(ctx, os.Getenv("DATABASE_URL"))
if err != nil {
slog.Error("db connect failed", "err", err)
os.Exit(1)
}
defer db.Close()
queue := pgqueue.NewQueue(db)
// HTTP-сервер: принимает загрузки и отдаёт метрики.
mux := http.NewServeMux()
mux.Handle("/metrics", promhttp.Handler())
mux.HandleFunc("/upload", func(w http.ResponseWriter, r *http.Request) {
// Упрощённо: энкодируем payload и кладём в очередь.
payload, _ := json.Marshal(ImagePayload{
S3Key: "uploads/" + r.URL.Query().Get("file"),
Widths: []int{320, 640, 1280},
Formats: []string{"webp", "jpg"},
})
if _, err := db.Exec(r.Context(),
`INSERT INTO jobs (type, payload) VALUES ('image_process', $1)`, payload); err != nil {
http.Error(w, "enqueue failed", http.StatusInternalServerError)
return
}
w.WriteHeader(http.StatusAccepted)
})
srv := &http.Server{Addr: ":8080", Handler: mux}
go srv.ListenAndServe()
// Запускаем пул воркеров.
go runWorkers(ctx, queue, 20)
<-ctx.Done()
slog.Info("shutting down")
shutCtx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
defer cancel()
srv.Shutdown(shutCtx)
}
Этот пример демонстрирует полный цикл: HTTP → PostgreSQL-очередь → worker pool с семафором → Prometheus-метрики → graceful shutdown. Семафор (chan struct{}) здесь заменяет явный пул горутин и является идиоматическим Go-паттерном для ограничения параллелизма.
Заключение
Worker pool в Go — не просто паттерн оптимизации, а фундаментальный инструмент построения надёжных высоконагруженных сервисов. Мы рассмотрели полный спектр: от базовой реализации на каналах до динамического масштабирования, retry без утечек горутин, graceful shutdown и PostgreSQL в роли персистентной очереди с FOR UPDATE SKIP LOCKED.
Ключевые выводы:
- Всегда ограничивайте параллелизм явно — через размер пула, семафор или буферизованный канал.
- Контекст должен пронизывать весь стек: от HTTP-хендлера до обращения к базе данных.
- Graceful shutdown с таймаутом обязателен в production — OS ждёт не бесконечно.
- Для персистентных задач PostgreSQL с
FOR UPDATE SKIP LOCKED— зрелое и надёжное решение, особенно в связке с библиотекой river. - Метрики Prometheus должны быть встроены с первого дня — не добавляйте их ретроспективно.
Go предоставляет все примитивы для реализации production-ready worker pool без внешних зависимостей. Понимание этих паттернов отличает разработчика, способного строить системы, устойчивые к реальным нагрузкам.
Технологии
Теги
Руслан Исмаилов
Senior Web / Backend разработчик. Senior web/backend разработчик с 9-летним опытом. Стек: PHP, Laravel, PostgreSQL, Redis, Docker, Kubernetes, REST, микросервисы, CI/CD. Подробнее обо мне →