Паттерны конкурентности: пул, producer-consumer, future, pipeline
В обзоре трека мы договорились: паттерн — это имя для решения, которое возникает снова и снова под давлением одних и тех же сил. В порождающих, структурных и поведенческих паттернах эти силы были про изменчивость: что-то в требованиях меняется чаще остального, и мы вставляем шов в нужном месте.
Здесь силы другие. Конкурентные паттерны выросли не из желания красиво расширять код, а из трёх физических ограничений, которые нельзя переспорить:
- Создание потока стоит дорого — примерно 1 МБ виртуальной памяти под стек в Linux/glibc и
десятки микросекунд на
clone(). Один поток на запрос при 50 000 RPS — это не «медленно», это «невозможно». - Разделяемое изменяемое состояние ломается — гонки, невидимость записей между ядрами, взаимные блокировки. Это не баг конкретного разработчика, а свойство модели.
- Скорости частей системы не совпадают — сеть отдаёт 10 000 сообщений в секунду, база переваривает 600. Разницу кто-то должен буферизовать или гасить.
Каждый паттерн ниже — ответ на одну из этих сил. Отсюда и главное правило чтения главы: не «какие бывают конкурентные паттерны», а «какая сила у меня в задаче преобладает».
Карта территории
конкурентности)) Управление ресурсом Thread Pool / Worker Pool work stealing размер пула по формуле Object Pool соединения к БД Half-Sync/Half-Async событийный слой + пул воркеров Передача работы Producer-Consumer ограниченная очередь обратное давление Pipeline стадии + каналы fan-out / fan-in Active Object очередь запросов к объекту Управление результатом Future / Promise композиция then/all отмена и таймаут Structured Concurrency время жизни = область видимости Completion Token коррелятор ответа Защита состояния Monitor Object мьютекс + условная переменная Immutability / Copy-on-Write Thread-Local Storage Read-Write Lock
Классический каталог здесь — не GoF, а POSA vol. 2 («Patterns for Concurrent and Networked Objects», Шмидт и др., 2000): именно там формализованы Active Object, Monitor Object, Reactor, Proactor, Half-Sync/Half-Async, Leader/Followers. Практическая настольная книга — «Java Concurrency in Practice» Брайана Гётца: она про Java, но модель рассуждений в ней универсальная.
Сила №1: поток дорог → Thread Pool
Интуиция
Представьте кол-центр. Наивная схема: на каждый звонок нанимаем нового оператора, а после звонка увольняем. Абсурд — найм длится дольше разговора. Реальная схема: держим 20 операторов, звонки кладём в общую очередь, освободившийся берёт следующий. Это и есть пул: отделяем время жизни исполнителя от времени жизни задачи.
Ровно это делает ExecutorService в Java, concurrent.futures.ThreadPoolExecutor в Python,
Task поверх ThreadPool в .NET. Go формально пул не даёт — там горутина стоит ~2–8 КБ и
мультиплексируется рантаймом, — но пул воркеров всё равно пишут вручную, чтобы ограничить
параллелизм (см. ниже про базу данных).
Жизненный цикл задачи в пуле
Ключ к пониманию пула — состояние Rejected. Пул без него (то есть с неограниченной очередью)
не пул, а бомба замедленного действия: задачи копятся, память кончается. CallerRunsPolicy —
элегантнейший трюк из Java: если очередь переполнена, задачу выполняет поток, который её
подал. Он на это время перестаёт подавать новые — система сама себя тормозит.
Сколько потоков?
Это самая частая ошибка проектирования: количество берут «на глаз». Есть формула из «Java Concurrency in Practice» (§8.2):
$$N_{threads} = N_{cpu} \times U_{target} \times \left(1 + \frac{W}{C}\right)$$
где W — среднее время ожидания (I/O, блокировки), C — среднее время счёта, U_target —
желаемая загрузка CPU (0…1).
- Чистый CPU-bound (
W/C ≈ 0):N = N_cpu(иногдаN_cpu + 1, чтобы закрыть page fault). - Типичный веб-хендлер: 5 мс счёта, 45 мс ожидания базы, 8 ядер,
U = 0.9:N = 8 × 0.9 × (1 + 45/5) = 72потока. - Полностью I/O-bound с ожиданием секундами — формула даёт тысячи потоков; это сигнал, что нужен не пул потоков, а асинхронная модель (см. раздел про future).
Но у формулы есть невидимая граница. Пул из 72 потоков, каждый из которых берёт соединение к
PostgreSQL, где max_connections = 100 и всего 4 диска, — это перегруз базы, а не ускорение.
Реальный предел ставит самый узкий ресурс за пулом, и его надо мерить, а не считать.
Почему нельзя грузить пул «под завязку»
Это следствие теории массового обслуживания. Для простейшей модели M/M/1 среднее время в системе
$$W = \frac{S}{1 - \rho}, \qquad \rho = \frac{\lambda}{\mu N}$$
При загрузке 50% задержка удваивается относительно чистой обработки. При 90% — вырастает в 10 раз. При 99% — в 100 раз. Отсюда практическое правило: проектируйте пул на 70–80% загрузки в пике, остальное — запас на всплеск. Инженеры, которые «оптимизировали» утилизацию до 95%, обычно потом объясняют, почему p99 улетел в потолок при том же среднем трафике.
Полезная спутница — закон Литтла L = λ × W: среднее число задач в системе равно частоте
поступления, умноженной на время в системе. Он даёт мгновенную оценку размера очереди без
всякого профилирования: 900 запросов/с × 80 мс = 72 запроса «в полёте» одновременно.
Код: пул с ограниченной очередью и отменой (Python)
import concurrent.futures as cf
import queue
import threading
import time
from dataclasses import dataclass
@dataclass
class Job:
job_id: int
payload: str
class BoundedPool:
"""Пул воркеров с ОГРАНИЧЕННОЙ очередью и корректной остановкой.
Стандартный ThreadPoolExecutor в Python использует неограниченную очередь:
submit() никогда не блокирует, и при перегрузе память течёт до OOM.
Здесь очередь ограничена, поэтому submit() создаёт обратное давление.
"""
_STOP = object() # sentinel: сигнал воркеру завершиться
def __init__(self, workers: int, capacity: int) -> None:
self._q: queue.Queue = queue.Queue(maxsize=capacity)
self._threads = [
threading.Thread(target=self._loop, name=f"worker-{i}", daemon=False)
for i in range(workers)
]
self._stopping = threading.Event()
for t in self._threads:
t.start()
def _loop(self) -> None:
while True:
item = self._q.get()
try:
if item is self._STOP:
return
self._handle(item)
except Exception: # noqa: BLE001
# КРИТИЧНО: исключение, вылетевшее из _loop, убивает воркер
# молча — пул деградирует до нуля потоков без единой строчки в логе.
import logging
logging.exception("job failed: %r", item)
finally:
self._q.task_done()
def _handle(self, job: Job) -> None:
time.sleep(0.01) # имитация работы
print(f"{threading.current_thread().name} -> {job.job_id}")
def submit(self, job: Job, timeout: float | None = None) -> bool:
"""Блокирует, если очередь полна. Возвращает False по таймауту."""
if self._stopping.is_set():
raise RuntimeError("pool is shutting down")
try:
self._q.put(job, timeout=timeout)
return True
except queue.Full:
return False # вызывающий решает: ждать, дропнуть, ретрай
def shutdown(self, drain: bool = True) -> None:
self._stopping.set()
if drain:
self._q.join() # дождаться выполнения принятых задач
for _ in self._threads:
self._q.put(self._STOP)
for t in self._threads:
t.join()
if __name__ == "__main__":
pool = BoundedPool(workers=4, capacity=16)
for i in range(50):
pool.submit(Job(i, f"payload-{i}"), timeout=1.0)
pool.shutdown()
Сложность. submit/take — амортизированные O(1) при очереди на связном списке или
кольцевом буфере. Память — O(capacity) на очередь плюс O(N × stack_size) на потоки; именно
второе слагаемое обычно и доминирует (72 потока × 1 МБ ≈ 72 МБ виртуального адресного
пространства, из них резидентно, как правило, десятки-сотни КБ на поток).
Contention. Одна общая очередь на 64 воркера — точка конкуренции: все дерутся за один мьютекс.
Отсюда work stealing (по одной deque на воркер, свободный крадёт с чужого хвоста):
ForkJoinPool в Java, Rayon в Rust, планировщик горутин в Go. Выигрыш заметен на десятках ядер
и мелких задачах; на 4 ядрах и задачах по 10 мс разницы вы не увидите.
Типичные ошибки с пулом
| Ошибка | Что происходит | Как чинить |
|---|---|---|
| Неограниченная очередь | Память растёт линейно, OOM через часы | maxsize, политика отказа |
| Исключение не ловится в цикле воркера | Пул тихо теряет потоки | try/except вокруг тела задачи |
| Задача в пуле ждёт результат другой задачи того же пула | Взаимоблокировка (thread starvation deadlock) | Разные пулы для разных уровней |
| Блокирующий вызов в пуле event loop | Встают все запросы, не только один | Отдельный пул: run_in_executor, Dispatchers.IO |
ThreadLocal не очищается |
Утечка на живущих вечно потоках | remove() в finally |
| Пул не останавливают | JVM/процесс не завершается (non-daemon) | shutdown() + awaitTermination |
Про thread starvation deadlock стоит сказать отдельно: пул из 10 потоков, каждая задача которого
внутри отправляет подзадачу в этот же пул и делает future.get(), встаёт намертво при 10
одновременных задачах. Все ждут исполнителей, которых больше нет. Это не гипотетика — типовой
инцидент, попавший в главу 8 «Java Concurrency in Practice».
Сила №3: скорости не совпадают → Producer-Consumer
Интуиция и связь с очередью
Producer-consumer — это Thread Pool, вывернутый наизнанку: там мы смотрели на исполнителей, здесь — на очередь между двумя подсистемами с разной скоростью. Очередь делает три вещи одновременно: развязывает по времени, развязывает по коду (производитель не знает потребителя) и — если ограничена — регулирует скорость.
Именно третий пункт большинство и упускает. queue.Queue() без maxsize, chan T без буфера
против make(chan T, 1000000), Kafka-consumer без лимита на in-flight — все они превращают
«медленный потребитель» в «упавший сервис». Ограниченная очередь переводит отказ из класса
«смерть по памяти» в класс «управляемая деградация»: производитель блокируется, ждёт, и система
сама выравнивается на скорости μ.
Обратное давление (backpressure) — это и есть свойство, ради которого стоит выбирать ограниченную очередь по умолчанию. Варианты реакции на переполнение, от мягкой к жёсткой:
- Блокировать производителя — идеально внутри процесса, опасно, если производитель обслуживает HTTP-запросы (клиенты копятся уже снаружи).
- Выполнить в потоке производителя (
CallerRunsPolicy) — самозамедление без блокировки. - Отбросить новое (load shedding) — правильно для метрик, телеметрии, «свежих» котировок.
- Отбросить старое — правильно, когда важно последнее значение (позиция курсора, снапшот).
- Вернуть ошибку 429/503 наверх — правильно на границе сервиса.
Выбор диктуется семантикой данных, а не вкусом. Отбрасывать платёжные события — катастрофа, отбрасывать точки графика загрузки CPU — норма.
Строгая формулировка: Monitor Object
Классическая реализация ограниченной очереди — паттерн Monitor Object: мьютекс + две
условные переменные (not_full, not_empty). Это тот случай, когда критично помнить два правила:
- ждать всегда в цикле
while, а не вif— из-за spurious wakeup и из-за того, что между сигналом и захватом мьютекса состояние мог изменить кто-то третий; - будить
notify_all, если ждущие ждут разных условий, иначе можно разбудить «не того» и потерять сигнал (классический баг lost wakeup).
import threading
from collections import deque
class BoundedBuffer:
"""Monitor Object: инвариант «0 <= len(items) <= capacity» защищён мьютексом."""
def __init__(self, capacity: int) -> None:
assert capacity > 0
self._capacity = capacity
self._items: deque = deque()
self._lock = threading.Lock()
# Две условные переменные на ОДНОМ мьютексе: producers и consumers
# ждут разных событий, поэтому notify адресный и lost wakeup не возникает.
self._not_full = threading.Condition(self._lock)
self._not_empty = threading.Condition(self._lock)
def put(self, item) -> None:
with self._not_full:
while len(self._items) == self._capacity: # именно while, не if
self._not_full.wait()
self._items.append(item)
self._not_empty.notify()
def take(self):
with self._not_empty:
while not self._items:
self._not_empty.wait()
item = self._items.popleft()
self._not_full.notify()
return item
В проде свою BoundedBuffer писать не надо — есть queue.Queue, ArrayBlockingQueue,
LinkedBlockingQueue, Channel в Kotlin, chan в Go. Но понимать её устройство надо: без этого
вы не отладите ни зависшего потребителя, ни «очередь не пустеет, хотя воркеры живые».
Go: тот же паттерн без единого мьютекса
package main
import (
"context"
"fmt"
"sync"
"time"
)
// Буферизованный канал — это и есть ограниченная очередь producer-consumer.
// Ёмкость 64 задаёт величину обратного давления: на 65-й записи producer встанет.
func main() {
ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second)
defer cancel()
jobs := make(chan int, 64)
// Производитель: закрывает канал — это единственный корректный сигнал «данных больше нет».
go func() {
defer close(jobs)
for i := 0; i < 1000; i++ {
select {
case jobs <- i:
case <-ctx.Done(): // отмена доходит и до заблокированной записи
return
}
}
}()
// Потребители: пул из 8 горутин на одном канале.
var wg sync.WaitGroup
for w := 0; w < 8; w++ {
wg.Add(1)
go func(id int) {
defer wg.Done()
for j := range jobs { // range по каналу завершится сам после close
time.Sleep(time.Millisecond)
_ = j
}
fmt.Printf("worker %d done\n", id)
}(w)
}
wg.Wait()
}
Два правила Go, нарушение которых даёт 90% зависаний: канал закрывает только отправитель
(закрытие потребителем даст панику у другого отправителя), и закрывает ровно один раз.
Если отправителей несколько, ставьте перед ними sync.WaitGroup и закрывайте канал в отдельной
горутине после wg.Wait().
Сила №2: результат приходит позже → Future / Promise
Интуиция
Future — это гардеробный номерок. Вы сдали пальто и получили не пальто, а обещание пальто. Номерок можно передать другому, положить в карман, обменять на пальто позже — а можно и не менять, если передумали. Разделение «начать работу» и «получить результат» — вся суть паттерна.
Терминология (её путают чаще всего): Future — сторона чтения, только получает результат.
Promise — сторона записи, тот, кто результат кладёт. В JavaScript слово Promise обозначает
оба конца сразу, в Java есть Future (чтение) и CompletableFuture (чтение + запись), в C++
разделение каноническое: std::promise / std::future.
продолжает работу P->>D: HTTP-запрос C->>F: then(map_result) — регистрирует колбэк D-->>P: ответ (через 300 мс) P->>F: complete(value) → состояние FULFILLED F->>C: колбэк вызывается на потоке-завершителе C->>F: get(timeout=1s) F-->>C: value (уже готово, возврат мгновенный) Note over C,D: альтернативная ветка D--xP: таймаут / 500 P->>F: completeExceptionally(err) → REJECTED F->>C: get() бросает исключение
Композиция — то, ради чего это делалось
Один future ценности почти не даёт: future.get() сразу после submit() — это обычный
синхронный вызов с накладными расходами. Ценность появляется, когда futures комбинируют:
import asyncio
import random
async def fetch_price(symbol: str) -> float:
await asyncio.sleep(random.uniform(0.05, 0.3)) # имитация сетевого вызова
if symbol == "BAD":
raise ValueError(f"нет котировки для {symbol}")
return round(random.uniform(10, 500), 2)
async def portfolio_value(symbols: list[str]) -> dict[str, float | str]:
"""Все запросы стартуют одновременно; общее время = max, а не sum."""
async with asyncio.TaskGroup() as tg: # Python 3.11+: структурная конкурентность
tasks = {s: tg.create_task(fetch_price(s)) for s in symbols}
return {s: t.result() for s, t in tasks.items()}
async def portfolio_tolerant(symbols: list[str]) -> dict[str, float | str]:
"""Вариант «частичный успех»: одна ошибка не убивает всю выборку."""
results = await asyncio.gather(
*(fetch_price(s) for s in symbols), return_exceptions=True
)
return {
s: (f"ERR: {r}" if isinstance(r, Exception) else r)
for s, r in zip(symbols, results)
}
async def with_deadline(symbol: str) -> float | None:
"""Таймаут — обязательная часть любого сетевого future, а не опция."""
try:
async with asyncio.timeout(0.2):
return await fetch_price(symbol)
except TimeoutError:
return None
if __name__ == "__main__":
print(asyncio.run(portfolio_tolerant(["AAPL", "MSFT", "BAD", "NVDA"])))
Базовые комбинаторы, которые есть во всех экосистемах под разными именами:
| Смысл | Python | JS | Java | Go |
|---|---|---|---|---|
| Преобразовать результат | await + выражение |
.then |
thenApply |
обычный код после <-ch |
| Цепочка асинхронных шагов | await подряд |
.then возвращает промис |
thenCompose |
последовательные вызовы |
| Все параллельно, ждём всех | gather / TaskGroup |
Promise.all |
allOf |
errgroup.Wait |
| Первый успешный | — | Promise.any |
— | select + отмена |
| Первый любой (гонка) | wait(FIRST_COMPLETED) |
Promise.race |
anyOf |
select |
| Обработка ошибки | try/except |
.catch |
exceptionally |
проверка err |
Обратите внимание на строку «цепочка»: thenApply против thenCompose — это ровно
разница между map и flatMap. Future — это функтор и монада, и об этом подробно будет в статье
«Функциональные паттерны». Если вы вернули
CompletableFuture<CompletableFuture<T>> — вы забыли flatMap.
Структурная конкурентность: главный сдвиг последних лет
Классическая проблема future: он живёт дольше того, кто его создал. Задача продолжает жечь CPU после того, как ответ клиенту уже отдан; исключение в ней некому обработать; отмена не доходит. «Go statement considered harmful» Мартина Сустрика/Натаниэля Смита (статья Smith, 2018) предложила правило: фоновая задача не может пережить блок, в котором создана.
Реализации: nursery в Trio, asyncio.TaskGroup (Python 3.11), StructuredTaskScope (Java 21+,
JEP 453), errgroup + context в Go, coroutineScope в Kotlin. Эффект — тот же, что у RAII
для памяти: время жизни задачи привязано к области видимости, а отмена и ошибки распространяются
автоматически.
// Go: errgroup — структурная конкурентность де-факто.
g, ctx := errgroup.WithContext(ctx)
g.SetLimit(8) // ограничение параллелизма прямо здесь
for _, url := range urls {
url := url
g.Go(func() error {
return fetch(ctx, url) // первая ошибка отменяет ctx для остальных
})
}
if err := g.Wait(); err != nil { // Wait не вернётся, пока живы дочерние горутины
return err
}
Ошибки с futures
get()сразу послеsubmit()— конкурентности нет, есть только оверхед.- Future без таймаута — один зависший сокет держит поток вечно; всегда
get(timeout)илиasyncio.timeout. - Проглоченное исключение — future, чей результат никто не прочитал, уносит ошибку с собой.
В Python
asyncioпро это хотя бы предупреждает («Task exception was never retrieved»), в JavaCompletableFutureмолчит. - Блокирующий вызов внутри асинхронного колбэка —
requests.get()в корутине встаёт весь event loop. Синхронный код выносится вrun_in_executor/asyncio.to_thread. - Смешение пулов — колбэк выполняется на потоке того, кто завершил future. Если это I/O-поток Netty, а вы в колбэке лезете в базу, вы блокируете сетевой цикл всего процесса.
Pipeline: конкурентность как композиция стадий
Интуиция
Конвейер на заводе: каждая станция делает одну операцию и передаёт деталь дальше. Пропускная способность конвейера равна пропускной способности самой медленной станции, а не сумме. Латентность одной детали при этом не улучшается — улучшается throughput.
Программный pipeline — цепочка стадий, соединённых очередями; каждая стадия может иметь свой параллелизм. Это композиция producer-consumer: выход стадии N — вход стадии N+1.
файлы/Kafka)] --> P1 subgraph St1["Стадия 1 · parse · CPU-bound"] P1[worker ×4] end P1 -->|chan cap=100| St2 subgraph St2["Стадия 2 · enrich · I/O-bound"] E1[worker ×32] end E2{{"fan-in"}} St2 --> E2 E2 -->|chan cap=500| St3 subgraph St3["Стадия 3 · batch+write · узкое место"] W1[worker ×2
батчи по 500] end St3 --> DB[(ClickHouse)] St3 -.->|ошибки| DLQ[(Dead letter queue)] style St3 fill:#e0704f22,stroke:#e0704f style E2 fill:#4fae7e22,stroke:#4fae7e
Здесь видно главное решение проектировщика: параллелизм каждой стадии подбирается под её
природу. CPU-стадия — по числу ядер. I/O-стадия — по формуле с W/C, десятки воркеров.
Стадия записи — по тому, сколько выдержит база (2 воркера с батчами по 500 строк почти всегда
быстрее 32 воркеров с одиночными вставками). Одинаковый параллелизм на всех стадиях —
самая частая ошибка новичка в потоковой обработке.
Fan-out / fan-in на Go
// generator — источник: превращает срез в канал (стадия 0).
func generator(ctx context.Context, nums ...int) <-chan int {
out := make(chan int)
go func() {
defer close(out)
for _, n := range nums {
select {
case out <- n:
case <-ctx.Done():
return
}
}
}()
return out
}
// stage — одна стадия обработки. Владелец выходного канала — сама стадия.
func stage(ctx context.Context, in <-chan int, f func(int) int) <-chan int {
out := make(chan int)
go func() {
defer close(out)
for v := range in {
select {
case out <- f(v):
case <-ctx.Done():
return
}
}
}()
return out
}
// fanIn — слияние N каналов в один: классический merge на WaitGroup.
func fanIn(ctx context.Context, chans ...<-chan int) <-chan int {
var wg sync.WaitGroup
out := make(chan int)
for _, c := range chans {
wg.Add(1)
go func(c <-chan int) {
defer wg.Done()
for v := range c {
select {
case out <- v:
case <-ctx.Done():
return
}
}
}(c)
}
go func() { wg.Wait(); close(out) }() // закрываем ПОСЛЕ всех отправителей
return out
}
func Run(ctx context.Context) int {
src := generator(ctx, 1, 2, 3, 4, 5, 6, 7, 8)
// fan-out: 4 копии одной стадии читают из общего входного канала
workers := make([]<-chan int, 4)
for i := range workers {
workers[i] = stage(ctx, src, func(n int) int { return n * n })
}
sum := 0
for v := range fanIn(ctx, workers...) { // fan-in обратно в один поток
sum += v
}
return sum
}
Анализ. Пусть стадии имеют пропускные способности μ₁ … μₖ. Тогда throughput всего
конвейера = min(μᵢ), латентность одного элемента = Σ (обработка + ожидание в очереди).
Память — O(Σ capacityᵢ × размер элемента); вот почему буферы стадий надо считать, а не
ставить «на глаз миллион». Ускорение по сравнению с последовательной обработкой ограничено
законом Амдала: если 10% работы принципиально последовательны, потолок ускорения — ×10,
сколько ядер ни добавляй. А универсальный закон масштабируемости Гуннера
(Neil Gunther, USL) добавляет
член когерентности: с ростом числа воркеров производительность не выходит на плато, а
начинает падать — из-за синхронизации между ними.
Где pipeline уже реализован за вас
- Go: Go blog, «Go Concurrency Patterns: Pipelines and cancellation» — каноническая статья.
- .NET: TPL Dataflow (
BufferBlock,TransformBlock,ActionBlockсBoundedCapacity). - Java: Reactive Streams / Project Reactor — там обратное давление часть спецификации (reactive-streams.org).
- Elixir/Erlang:
GenStageиFlow— потребитель явно запрашиваетdemand, то есть backpressure тянущий, а не толкающий (см. обзор Elixir). - Python:
asyncio.Queueмежду корутинами; для тяжёлого CPU —multiprocessingили Ray.
Как выбрать паттерн
Практическая процедура, а не список:
- Измерьте
W/C. Если ожидание не доминирует — вам, возможно, нужен не параллелизм, а оптимизация алгоритма. УскоритьO(n²)доO(n log n)дешевле, чем распараллелить на 8 ядер. - Найдите узкое место. Оно почти всегда одно и почти всегда снаружи процесса (база, диск, внешний API). Параллелизм перед узким местом только увеличит очередь.
- Начните с пула и ограниченной очереди. Это покрывает большинство задач и понятно любому в команде.
- Асинхронность — когда счёт соединений идёт на тысячи. 10 000 одновременных WebSocket на пуле потоков не живут; на event loop — живут.
- Актор/pipeline — когда есть естественные стадии или изолированное состояние. Не раньше.
- Lock-free — почти никогда. Правильная реализация неблокирующей структуры данных — это недели работы и модель памяти в голове. Возьмите готовую из стандартной библиотеки.
Что ломается в проде: короткий разбор классики
Пул перед базой без лимита соединений. 200 воркеров, max_connections=100. Половина
воркеров вечно ждёт соединение, метрики показывают «пул занят», CPU простаивает. Лечение —
пул соединений (HikariCP, pgbouncer) меньше пула потоков и явный таймаут получения соединения.
Очередь как «буфер на всякий случай». Kafka-consumer читает быстрее, чем пишет в БД, но
разработчик ставит внутреннюю очередь на 1 000 000 сообщений «чтобы не терять». В инциденте
эта очередь заполняется за 15 минут, процесс падает по OOM, offset не закоммичен, после
рестарта всё повторяется — вечный цикл. Лечение: очередь на 1 000, обратное давление до самого
poll().
Отмена, которая не отменяет. future.cancel(true) в Java не убивает поток — он лишь ставит
флаг прерывания. Задача, которая не проверяет Thread.interrupted() и не ловит
InterruptedException, продолжает работать. То же с ctx.Done() в Go: контекст — это
уведомление, а не kill. Правило: любой длинный цикл проверяет отмену.
Ложная конкурентность. Каждый запрос уходит в пул, но внутри берёт один и тот же
synchronized-объект (например, кэш на HashMap под глобальным замком). Реального
параллелизма нет, есть только оверхед на переключение контекста и очередь на мьютексе. Профайлер
покажет «CPU 5%, latency 800 мс» — характерная подпись.
Метрики, которых нет. Минимальный набор для любой конкурентной подсистемы: длина очереди (gauge), время ожидания в очереди (histogram), число активных воркеров, счётчик отказов/дропов, счётчик исключений в воркерах. Без длины очереди вы не отличите «медленно работает» от «захлебнулись»; это два разных инцидента с разным лечением.
Мини-итог
- Конкурентные паттерны решают три силы: дорогой поток, разделяемое состояние, несовпадение скоростей. Определите свою — паттерн выберется почти автоматически.
- Thread Pool отделяет жизнь исполнителя от жизни задачи; размер считается по
N_cpu × U × (1 + W/C), а рабочая загрузка держится на 70–80% из-за колена кривойS/(1−ρ). - Producer-Consumer развязывает подсистемы, но ценность даёт только ограниченная очередь: она превращает отказ по памяти в управляемую деградацию.
- Future/Promise отделяет запуск от получения результата; настоящая польза — в композиции
(
all,race,flatMap) и в структурной конкурентности, где задача не переживает свою область. - Pipeline — композиция стадий, где throughput равен минимуму по стадиям, а параллелизм каждой стадии подбирается под её природу (CPU / I/O / внешний лимит).
- Общий знаменатель всех четырёх: явные границы. Ограниченный пул, ограниченная очередь, таймаут на future, ёмкость канала. Всё безграничное в конкурентном коде рано или поздно съедает память или ядро.
Источники
- Brian Goetz et al., «Java Concurrency in Practice», 2006 — jcip.net; главы 6–8 (пулы) и 5 (ограниченные очереди) — обязательное чтение независимо от языка.
- Douglas Schmidt et al., «Pattern-Oriented Software Architecture, Vol. 2: Patterns for Concurrent and Networked Objects», 2000 — сайт автора.
- «Go Concurrency Patterns: Pipelines and cancellation» — go.dev/blog/pipelines.
- Nathaniel J. Smith, «Notes on structured concurrency, or: Go statement considered harmful», 2018 — vorpus.org.
- JEP 444 «Virtual Threads» и JEP 453 «Structured Concurrency» — openjdk.org/jeps/444, openjdk.org/jeps/453.
- Reactive Streams specification (обратное давление как контракт) — reactive-streams.org.
- Neil Gunther, Universal Scalability Law — perfdynamics.com.
- Документация: Python
concurrent.futures, Pythonasyncio.TaskGroup, Gosync/errgroup, TPL Dataflow, Elixir GenStage.
Что дальше
Мы разобрали паттерны, которые управляют временем и разделяемым состоянием. Следующий шаг — паттерны, которые от разделяемого состояния избавляются: неизменяемые данные и функции как единица композиции. Это отдельная традиция со своим словарём.
Функциональные паттерны: функторы, монады, линзы, комбинаторы