6.11. Worker processes
Цели
После этого материала вы сможете:
- рассчитать число worker-процессов исходя из лимита container, а не из числа CPU хоста;
- объяснить, почему
os.cpu_count()даёт неверный ответ при--cpus; - выбрать между процессами и потоками для вашего типа нагрузки;
- оценить расход памяти на worker и понять, что даёт
preload; - написать background worker, который корректно завершается по
SIGTERMи не теряет задачи; - объяснить, почему обработчик задач обязан быть идемпотентным;
- организовать периодические задачи без
cronвнутри container.
Предварительные знания
- 2.4. cgroups — механизм лимитов;
- 6.5. Сигналы и PID 1 в Python;
- 6.9. Flask и 6.10. FastAPI — серверы приложений.
Рабочий пример — resources/examples/worker-redis/.
Ключевые термины
| Термин | Объяснение |
|---|---|
worker | Процесс, обрабатывающий запросы или задачи |
pre-fork | Модель: master создаёт worker'ов через fork до приёма нагрузки |
preload | Загрузка приложения в master до fork |
copy-on-write | Механизм ядра: после fork страницы памяти общие, копируются при записи |
at-least-once | Гарантия: задача выполнится минимум один раз, возможно — больше |
идемпотентность | Свойство: повторное выполнение даёт тот же результат |
cpu.max | Файл cgroup v2 с квотой CPU |
Теория
Формула, которая не работает в container
Документация Gunicorn предлагает исходную точку:
Generally we recommend
(2 x $num_cores) + 1as the number of workers to start off with.
Формула разумна на выделенном сервере. В container она даёт неверный результат, потому что «num_cores» обычно получают так:
workers = 2 * os.cpu_count() + 1
os.cpu_count() возвращает число CPU хоста. Container с --cpus 1.5 на 16-ядерной машине получит 2 * 16 + 1 = 33 worker'а при квоте в полтора ядра.
Последствия конкретны:
| Что происходит | Почему |
|---|---|
| Память × 33 | Каждый worker — отдельный интерпретатор с копией импортов |
OOM kill (код 137) | Суммарный RSS превышает --memory |
| Рост latency | 33 процесса делят квоту 1.5 CPU; каждый получает срез |
| Больше переключений контекста | Планировщик перебирает процессы, которым не хватает квоты |
Три разных числа
Внутри container «сколько у меня CPU» имеет три разных ответа.
| Способ | Что возвращает | Учитывает --cpus | Учитывает --cpuset-cpus |
|---|---|---|---|
os.cpu_count() | CPU, видимые ядру | нет | нет |
len(os.sched_getaffinity(0)) | CPU в маске affinity | нет | да |
/sys/fs/cgroup/cpu.max | Квота cgroup | да | нет |
Причина в природе механизмов (урок 2.4):
--cpuset-cpusменяет маску affinity — процесс физически не запускается на других ядрах, иsched_getaffinityэто видит;--cpusзадаёт квоту времени: процесс может исполняться на любом ядре, но суммарно не более заданной доли. Маска не меняется, поэтому ни один системный вызов о лимите не сообщает.
/proc/cpuinfo в container тоже показывает CPU хоста: procfs не виртуализирован по CPU.
Правильный расчёт читает cgroup:
def available_cpus() -> float:
"""Число CPU с учётом лимита container (cgroup v2)."""
try:
quota, period = Path("/sys/fs/cgroup/cpu.max").read_text().split()
except (OSError, ValueError):
return float(len(os.sched_getaffinity(0)))
if quota == "max": # лимит не задан
return float(len(os.sched_getaffinity(0)))
return int(quota) / int(period)
Формат cpu.max — два числа: квота и период в микросекундах. --cpus 1.5 даёт 150000 100000.
Проще: задать число явно
Расчёт по cgroup работает, но у него есть недостаток: число worker'ов становится неявным и меняется при изменении лимита. В production предпочтительнее противоположный подход — задать число явно переменной окружения:
services:
api:
environment:
WEB_CONCURRENCY: "3"
deploy:
resources:
limits:
cpus: "1.5"
memory: 512M
Переменную WEB_CONCURRENCY читает Gunicorn (документировано) — и, как показано ниже, Uvicorn тоже.
Преимущества явного задания:
| Свойство | Значение |
|---|---|
| Число видно в конфигурации | Не нужно вычислять, читая код |
| Меняется без пересборки образа | Подбор под нагрузку — вопрос перезапуска |
| Согласовано с лимитом памяти | Оба числа рядом, видно соотношение |
| Одинаково при любом способе запуска | Не зависит от того, что вернёт cgroup |
Правило: автоматический расчёт — для сред разработки, явное число — для production.
Процессы или потоки
Два способа обрабатывать несколько запросов одновременно, с разными свойствами.
Процессы (--workers) | Потоки (--threads) | |
|---|---|---|
| Память | Полная копия на каждый | Общая |
| GIL | У каждого свой | Один на процесс |
| CPU-bound нагрузка | Масштабируется | Не масштабируется |
| I/O-bound нагрузка | Масштабируется | Масштабируется |
| Падение одного | Остальные живут | Процесс падает целиком |
| Утечка памяти | Изолирована, лечится max_requests | Общая |
GIL определяет главное ограничение: два потока одного процесса не исполняют Python-байткод одновременно. Потоки помогают только тогда, когда время уходит на ожидание — сеть, диск, база.
Практическое соответствие типу нагрузки:
| Нагрузка | Конфигурация |
|---|---|
| CPU-bound (вычисления, сериализация больших ответов) | Только процессы: --workers N |
| I/O-bound синхронный (запросы к базе через синхронный драйвер) | --workers N --threads M |
| I/O-bound асинхронный (FastAPI) | Процессы; параллелизм даёт цикл событий |
| Смешанная | Процессы плюс немного потоков; измерять |
Для FastAPI потоки не нужны: цикл событий уже обслуживает тысячи ожидающих соединений внутри одного процесса (урок 6.10).
Память на worker
Каждый процесс — отдельный интерпретатор со своей копией импортированных модулей. Базовый расход:
| Компонент | Порядок |
|---|---|
| Интерпретатор Python | 10–15 MB |
| Flask плюс зависимости | +20–30 MB |
| FastAPI плюс Pydantic | +40–60 MB |
| Модели ML, кэши, данные | зависит от приложения |
Отсюда правило подбора лимита памяти:
--memory ≥ (RSS одного worker'а × число worker'ов) × 1.3
Запас в 30 % покрывает фрагментацию аллокатора и пики при обработке запроса (урок 6.13).
Что даёт preload
Без preload каждый worker импортирует приложение сам после fork. С preload master импортирует приложение один раз, а fork наследует уже загруженную память.
# gunicorn.conf.py
preload_app = True
Ожидаемая выгода — copy-on-write: после fork страницы физически общие, копируются только при записи. Значит, код и константы приложения существуют в одном экземпляре.
Но в CPython эта выгода частично теряется. Причина — счётчики ссылок: они хранятся в заголовке каждого объекта. Любое обращение к объекту меняет ob_refcnt, то есть пишет в его страницу, и ядро копирует страницу целиком. Сборщик мусора усугубляет: он обходит объекты и трогает их все.
Смягчение — gc.freeze(), доступный с Python 3.7:
# gunicorn.conf.py
import gc
preload_app = True
def when_ready(server):
"""Вызывается в master после загрузки приложения, до fork worker'ов."""
gc.freeze() # переносит объекты в «постоянное» поколение
gc.freeze() перемещает существующие объекты в поколение, которое сборщик не обходит. Это уменьшает число тронутых страниц, но не устраняет проблему полностью: обычные обращения к объектам всё равно меняют refcount.
Второе следствие preload_app, менее очевидное и более важное на практике:
Ресурсы, созданные до
fork, наследуются всеми worker'ами.
Соединение с базой, открытое при импорте, будет общим файловым дескриптором в нескольких процессах — что приводит к перемешанным ответам и повреждению протокола. Правило: соединения создавать после fork, в хуке post_fork или при первом использовании в worker'е.
Свойство preload_app | Оценка |
|---|---|
| Экономия памяти | Есть, но меньше ожидаемой из-за refcount |
| Скорость старта worker'ов | Заметно выше |
| Ошибки импорта видны сразу | Плюс: падает master, а не каждый worker по очереди |
Перезапуск по SIGHUP без остановки | Не работает: код загружен в master |
| Наследование соединений | Риск: требует создавать их после fork |
Background worker
Второй тип worker'а — не обслуживающий HTTP, а разбирающий очередь задач. Его требования к container другие.
| Свойство | Web worker | Background worker |
|---|---|---|
| Входящие соединения | Есть | Нет |
| Порт | Публикуется | Не нужен |
HEALTHCHECK | HTTP-запрос | Нет естественного endpoint |
| Единица работы | Запрос, миллисекунды | Задача, секунды или минуты |
| Grace period | 10 секунд обычно хватает | Часто нужно больше |
| Потеря при остановке | Один запрос | Целая задача |
Ключевое требование — цикл должен проверять флаг завершения:
_shutdown = False
def request_shutdown(signum, _frame):
global _shutdown
_shutdown = True # только флаг, никакой работы в обработчике
signal.signal(signal.SIGTERM, request_shutdown)
while not _shutdown:
task = queue.get(timeout=2) # таймаут обязателен
if task is None:
continue
process(task)
Два обязательных элемента:
- Обработчик только выставляет флаг. Обработчик сигнала прерывает основной поток в произвольной точке; сложная работа в нём небезопасна (урок 6.5).
- Ожидание задачи с таймаутом. Блокирующее ожидание без таймаута не даст циклу проверить флаг: worker будет висеть до
SIGKILL.
Grace period задаётся отдельно, так как задача длиннее запроса:
services:
worker:
stop_grace_period: 30s
Почему обработчик обязан быть идемпотентным
Очереди дают гарантию at-least-once, а не exactly-once. Задача может быть выполнена дважды. Сценарии:
| Сценарий | Результат |
|---|---|
Worker убит SIGKILL после обработки, но до подтверждения | Задача вернётся в очередь и выполнится снова |
| Сеть оборвалась при отправке подтверждения | То же |
| Задача поставлена повторно из-за ошибки на стороне отправителя | Дубль |
| Таймаут видимости истёк, задачу забрал другой worker | Два worker'а выполняют одну задачу |
Exactly-once в распределённой системе недостижим в общем случае — доказано теоретически. Практический ответ — сделать повтор безвредным:
| Приём | Как работает |
|---|---|
| Ключ идемпотентности | SETNX task:{id}:done перед выполнением; если ключ есть — пропустить |
| Естественная идемпотентность | UPDATE ... SET status='sent' вместо counter = counter + 1 |
| Upsert вместо insert | INSERT ... ON CONFLICT DO NOTHING |
| Проверка состояния | Перед отправкой письма проверить, не отправлено ли |
Надёжность очереди на Redis
Простейший способ забрать задачу — BLPOP: он удаляет элемент из списка. Если worker будет убит после BLPOP, но до завершения обработки, задача исчезнет.
Надёжный вариант — атомарно переместить задачу в список «в работе»:
BLMOVE tasks processing LEFT RIGHT 2
Задача одновременно покидает очередь и попадает в processing. После успешной обработки её оттуда удаляют (LREM). Задачи, застрявшие в processing, возвращает в очередь отдельный процесс-«уборщик».
| Команда | Атомарность | Потеря при падении worker'а |
|---|---|---|
BLPOP tasks 2 | Атомарна | Задача потеряна |
BLMOVE tasks processing LEFT RIGHT 2 | Атомарна | Задача остаётся в processing |
BRPOPLPUSH делает то же, но объявлен устаревшим начиная с Redis 6.2 в пользу BLMOVE.
Учебный пример курса использует BLPOP ради простоты; в разделе «Практическое упражнение» вы переведёте его на BLMOVE.
Для production обычно берут готовую библиотеку — Celery, RQ, Dramatiq, — которая решает это и десяток смежных вопросов. Понимать механизм всё равно необходимо: настройки этих библиотек описывают именно его.
Периодические задачи
cron внутри container — распространённый антипаттерн:
| Проблема | Следствие |
|---|---|
cron становится PID 1 | Приложение — его потомок; сигналы обрабатывает не оно |
cron пишет в syslog или почту | Логи не попадают в docker logs |
Собственное окружение cron | Переменные окружения container недоступны заданию |
| Часовой пояс по умолчанию UTC | Задание выполняется не тогда, когда ожидалось |
| Нет защиты от наложения | Долгое задание запускается повторно поверх предыдущего |
Варианты по возрастанию сложности инфраструктуры:
| Вариант | Когда подходит |
|---|---|
Планировщик оркестратора (Kubernetes CronJob, Swarm) | Есть оркестратор — лучший вариант |
systemd-таймер на хосте, запускающий docker run | Одиночный сервер, задача — отдельный процесс |
| Планировщик внутри приложения (APScheduler, Celery beat) | Задача тесно связана с кодом приложения |
Цикл while True: … sleep(interval) в отдельном container | Простые интервальные задачи |
Общее требование ко всем: защита от одновременного выполнения. При нескольких репликах планировщик запустится в каждой. Решение — распределённая блокировка:
SET cron:daily-report locked NX EX 3600
Только один экземпляр получит OK и выполнит задачу.
Внутренний механизм
Что происходит при fork
- Ядро создаёт копию процесса; таблицы страниц указывают на те же физические страницы, помеченные read-only.
- При записи в такую страницу возникает page fault.
- Ядро копирует страницу и снимает пометку.
Отсюда «copy-on-write»: копируется не всё, а только изменённое. В CPython изменяется больше, чем кажется, — из-за refcount.
Наследуются: открытые файловые дескрипторы, отображённая память, обработчики сигналов. Не наследуются: потоки (в дочернем остаётся только вызвавший fork).
Последнее — причина частой ошибки: пул соединений с фоновым потоком, созданный до fork, в worker'е окажется без своего потока и перестанет работать.
Почему --cpus невидим для процесса
--cpus 1.5 записывает в cgroup:
/sys/fs/cgroup/cpu.max → "150000 100000"
Планировщик ядра учитывает эту квоту, но интерфейс sched_getaffinity возвращает маску разрешённых CPU, а не квоту. Маска не менялась. Программе доступна только сама cgroup-файловая система — что и делает функция available_cpus() выше.
Команды и примеры
Три ответа на вопрос «сколько CPU»
cat > /tmp/cpus.py <<'PY'
"""Три способа узнать число CPU внутри container."""
import os
from pathlib import Path
def cgroup_cpus() -> str:
p = Path("/sys/fs/cgroup/cpu.max")
if not p.exists():
return "нет (cgroup v1?)"
quota, period = p.read_text().split()
if quota == "max":
return "лимит не задан"
return f"{int(quota) / int(period):.2f}"
print(f" os.cpu_count(): {os.cpu_count()}")
print(f" len(sched_getaffinity(0)): {len(os.sched_getaffinity(0))}")
print(f" cgroup cpu.max: {cgroup_cpus()}")
print(f" формула 2*cpu_count()+1: {2 * os.cpu_count() + 1} worker'ов")
PY
echo "═══ без лимита ═══"
docker run --rm -v /tmp/cpus.py:/c.py:ro python:3.13-slim python /c.py
echo "═══ --cpus 1.5 ═══"
docker run --rm --cpus 1.5 -v /tmp/cpus.py:/c.py:ro python:3.13-slim python /c.py
echo "═══ --cpuset-cpus 0,1 ═══"
docker run --rm --cpuset-cpus 0,1 -v /tmp/cpus.py:/c.py:ro python:3.13-slim python /c.py
Ожидаемый вывод на 8-ядерной машине:
═══ без лимита ═══
os.cpu_count(): 8
len(sched_getaffinity(0)): 8
cgroup cpu.max: лимит не задан
формула 2*cpu_count()+1: 17 worker'ов
═══ --cpus 1.5 ═══
os.cpu_count(): 8
len(sched_getaffinity(0)): 8
cgroup cpu.max: 1.50
формула 2*cpu_count()+1: 17 worker'ов
═══ --cpuset-cpus 0,1 ═══
os.cpu_count(): 2
len(sched_getaffinity(0)): 2
cgroup cpu.max: лимит не задан
формула 2*cpu_count()+1: 5 worker'ов
Главное наблюдение — средний блок: при квоте 1.5 CPU обычная формула предлагает 17 worker'ов, потому что оба Python-вызова видят 8 CPU хоста. Только cpu.max знает правду.
Обратите внимание и на различие в третьем блоке: --cpuset-cpus виден обоим вызовам, --cpus — ни одному.
Сколько памяти стоит worker
mkdir -p /tmp/wk && cd /tmp/wk
cat > app.py <<'PY'
import os
from flask import Flask
app = Flask(__name__)
@app.get("/")
def index():
return {"pid": os.getpid()}
PY
cat > requirements.txt <<'EOF'
flask==3.1.3
gunicorn==26.0.0
EOF
cat > Dockerfile <<'EOF'
FROM python:3.13-slim
ENV PYTHONUNBUFFERED=1
WORKDIR /app
COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt
COPY app.py .
CMD ["sh", "-c", "exec gunicorn -w ${WEB_CONCURRENCY:-1} -b 0.0.0.0:8000 app:app"]
EOF
docker build -q -t wk . > /dev/null
for n in 1 2 4 8; do
docker run -d --name "wk$n" -e WEB_CONCURRENCY="$n" wk > /dev/null
sleep 4
mem="$(docker stats --no-stream --format '{{.MemUsage}}' "wk$n")"
procs="$(docker exec "wk$n" sh -c 'ps -o pid= | wc -l')"
printf ' workers=%-2s процессов=%-2s память: %s\n' "$n" "$procs" "$mem"
docker rm -f "wk$n" > /dev/null
done
Ожидаемый вывод (порядок величины; конкретные числа зависят от машины):
workers=1 процессов=3 память: 42.1MiB / 15.6GiB
workers=2 процессов=4 память: 68.3MiB / 15.6GiB
workers=4 процессов=6 память: 121MiB / 15.6GiB
workers=8 процессов=10 память: 227MiB / 15.6GiB
Рост близок к линейному: примерно 26 MB на дополнительный worker плюс постоянные ~16 MB на master. Отсюда прямо считается лимит: при --memory 256M восемь worker'ов уже на грани.
Процессов на два больше числа worker'ов: master Gunicorn плюс сам ps.
Что происходит при завышенном числе worker'ов
echo "═══ workers=17 при --memory 256M и --cpus 1.5 ═══"
docker run -d --name wk-over --memory 256M --cpus 1.5 -e WEB_CONCURRENCY=17 wk > /dev/null
sleep 12
printf ' статус: %s\n' "$(docker inspect wk-over --format '{{.State.Status}}')"
printf ' OOMKilled: %s\n' "$(docker inspect wk-over --format '{{.State.OOMKilled}}')"
printf ' код выхода: %s\n' "$(docker inspect wk-over --format '{{.State.ExitCode}}')"
docker logs wk-over 2>&1 | tail -3
docker rm -f wk-over > /dev/null
Ожидаемый вывод:
═══ workers=17 при --memory 256M и --cpus 1.5 ═══
статус: exited
OOMKilled: true
код выхода: 137
[ERROR] Worker (pid:24) was sent SIGKILL! Perhaps out of memory?
Код 137 и OOMKilled: true — подпись OOM killer'а (урок 6.13). Именно этим заканчивается формула 2 * os.cpu_count() + 1 в container с лимитом.
Обратите внимание: Gunicorn сообщает о гибели worker'а, но причину указывает предположительно («Perhaps out of memory?»). Достоверный ответ даёт docker inspect.
Проверка WEB_CONCURRENCY в Uvicorn
Gunicorn читает эту переменную по документации. Проверим Uvicorn:
cd resources/examples/fastapi-basic
docker build -q -t fb . > /dev/null
docker run -d --name fbw -e WEB_CONCURRENCY=3 -p 8000:8000 fb > /dev/null
sleep 7
printf 'процессов python: %s\n' \
"$(docker exec fbw sh -c 'for f in /proc/[0-9]*/cmdline; do tr "\0" " " < "$f"; echo; done | grep -c "[p]ython"')"
docker rm -f fbw > /dev/null; docker rmi -f fb > /dev/null
Ожидаемый вывод:
процессов python: 4
Четыре процесса — менеджер и три worker'а: переменная учтена. Если у вас получилось 1, ваша версия Uvicorn её не читает — тогда задавайте --workers явно. Это ровно тот случай, когда стоит проверить командой, а не поверить тексту.
Потоки против процессов на CPU-bound нагрузке
cd /tmp/wk
cat > cpu_app.py <<'PY'
"""CPU-bound endpoint: демонстрация влияния GIL."""
import os
def burn(n: int = 220_000) -> int:
total = 0
for i in range(n):
total += i * i % 7
return total
def app(environ, start_response):
burn()
start_response("200 OK", [("Content-Type", "text/plain")])
return [f"{os.getpid()}\n".encode()]
PY
cat > Dockerfile.cpu <<'EOF'
FROM python:3.13-slim
ENV PYTHONUNBUFFERED=1
WORKDIR /app
RUN pip install --no-cache-dir gunicorn==26.0.0
COPY cpu_app.py .
CMD ["sh", "-c", "exec gunicorn ${GOPTS} -b 0.0.0.0:8000 cpu_app:app"]
EOF
docker build -q -f Dockerfile.cpu -t wkcpu . > /dev/null
bench() {
local name="$1" opts="$2"
docker run -d --name bench --cpus 2 -e GOPTS="$opts" -p 8001:8000 wkcpu > /dev/null
sleep 4
local start end
start="$(date +%s.%N)"
for _ in $(seq 8); do curl -s localhost:8001/ > /dev/null & done
wait
end="$(date +%s.%N)"
printf ' %-28s %.2f c\n' "$name" "$(awk -v a="$start" -v b="$end" 'BEGIN{print b-a}')"
docker rm -f bench > /dev/null
}
echo "═══ 8 параллельных CPU-bound запросов, --cpus 2 ═══"
bench "1 процесс" "-w 1"
bench "1 процесс, 4 потока" "-w 1 --threads 4"
bench "2 процесса" "-w 2"
bench "4 процесса" "-w 4"
Ожидаемый вывод (относительные значения важнее абсолютных):
═══ 8 параллельных CPU-bound запросов, --cpus 2 ═══
1 процесс 4.10 c
1 процесс, 4 потока 4.18 c
2 процесса 2.15 c
4 процесса 2.20 c
Три вывода из этого измерения:
- Потоки не помогли — 4.18 против 4.10 секунды. GIL не даёт двум потокам исполнять байткод одновременно; появились лишь накладные расходы на переключение.
- Процессы помогли вдвое — у каждого свой GIL, и оба использовали квоту 2 CPU.
- Четыре процесса не дали ничего сверх двух — квота исчерпана на двух. Дальнейшее увеличение только тратит память.
Последнее — практическое доказательство того, что число worker'ов должно исходить из лимита container.
docker rmi -f wk wkcpu > /dev/null 2>&1; rm -rf /tmp/wk /tmp/cpus.py
Background worker: graceful shutdown
cd resources/examples/worker-redis
docker compose up -d --build
sleep 12
docker compose ps --format 'table {{.Service}}\t{{.Status}}'
Ожидаемый вывод:
SERVICE STATUS
redis Up 12 seconds (healthy)
worker Up 6 seconds
Поставим задачи и остановим worker в момент обработки:
docker compose exec -T worker python producer.py 5 3
sleep 4
echo "═══ docker compose stop во время обработки ═══"
s="$(date +%s.%N)"
docker compose stop worker
e="$(date +%s.%N)"
printf ' время остановки: %.1f c\n' "$(awk -v a="$s" -v b="$e" 'BEGIN{print b-a}')"
printf ' код выхода: %s\n' "$(docker compose ps -a --format '{{.ExitCode}}' worker)"
docker compose logs worker --no-log-prefix | tail -6
Ожидаемый вывод:
поставлено задач: 5 (по 3.0 c каждая)
длина очереди: 5
═══ docker compose stop во время обработки ═══
время остановки: 1.9 c
код выхода: 0
2026-07-30 12:03:44 INFO [worker] задача 2: начата
2026-07-30 12:03:45 INFO [worker] получен SIGTERM — завершаю после текущей задачи
2026-07-30 12:03:47 INFO [worker] задача 2: завершена
2026-07-30 12:03:47 INFO [worker] остановлен штатно, обработано задач: 2
Разберём последовательность:
| Событие | Что показывает |
|---|---|
| Сигнал получен во время задачи 2 | Обработчик отработал немедленно |
| Задача 2 доведена до конца | Обработчик только выставил флаг, работу не прервал |
Выход с кодом 0 | Штатное завершение, не 137 |
| Остановка за 1.9 с, а не за 30 с | Не понадобился весь stop_grace_period |
Три задачи остались в очереди — это правильно: их заберёт следующий worker.
printf 'осталось в очереди: %s\n' \
"$(docker compose exec -T redis redis-cli llen tasks)"
Ожидаемый вывод:
осталось в очереди: 3
Почему BLPOP теряет задачу
Проверим, что происходит при жёстком убийстве:
docker compose up -d worker
sleep 4
docker compose exec -T worker python producer.py 3 8 > /dev/null
sleep 3
echo "═══ до SIGKILL ═══"
printf ' в очереди: %s\n' "$(docker compose exec -T redis redis-cli llen tasks)"
docker compose kill -s SIGKILL worker > /dev/null
sleep 1
echo "═══ после SIGKILL ═══"
printf ' в очереди: %s\n' "$(docker compose exec -T redis redis-cli llen tasks)"
printf ' код выхода: %s\n' "$(docker compose ps -a --format '{{.ExitCode}}' worker)"
echo " задача, которая обрабатывалась, нигде не осталась — потеряна"
Ожидаемый вывод:
═══ до SIGKILL ═══
в очереди: 2
═══ после SIGKILL ═══
в очереди: 2
код выхода: 137
задача, которая обрабатывалась, нигде не осталась — потеряна
Задач было три: одну worker забрал (BLPOP удалил её из списка), две остались. После SIGKILL в очереди по-прежнему две — забранная задача исчезла вместе с процессом.
Это не ошибка примера, а неизбежное свойство BLPOP. Устранение — в упражнении.
docker compose down -v
Периодическая задача с блокировкой
mkdir -p /tmp/sched && cd /tmp/sched
cat > tick.py <<'PY'
"""Периодическая задача с защитой от параллельного выполнения."""
from __future__ import annotations
import os
import signal
import sys
import time
import redis
INTERVAL = float(os.environ.get("INTERVAL", "3"))
LOCK_KEY = os.environ.get("LOCK_KEY", "cron:tick")
NAME = os.environ.get("REPLICA_NAME", "?")
_stop = False
def _handle(signum, _frame):
global _stop
_stop = True
def main() -> int:
signal.signal(signal.SIGTERM, _handle)
signal.signal(signal.SIGINT, _handle)
client = redis.Redis.from_url(
os.environ.get("REDIS_URL", "redis://redis:6379/0"),
decode_responses=True,
)
for _ in range(30):
try:
client.ping()
break
except redis.ConnectionError:
time.sleep(1)
while not _stop:
# NX: ключ ставится только если его нет. EX: срок жизни меньше интервала,
# иначе следующий запуск будет заблокирован предыдущим.
got = client.set(LOCK_KEY, NAME, nx=True, ex=int(INTERVAL) - 1)
if got:
print(f"[{NAME}] выполняю задачу", flush=True)
else:
print(f"[{NAME}] пропускаю: занято", flush=True)
# sleep дробится, чтобы SIGTERM не ждал полного интервала
slept = 0.0
while slept < INTERVAL and not _stop:
time.sleep(0.2)
slept += 0.2
print(f"[{NAME}] остановлен штатно", flush=True)
return 0
if __name__ == "__main__":
sys.exit(main())
PY
cat > requirements.txt <<'EOF'
redis==8.1.0
EOF
cat > Dockerfile <<'EOF'
FROM python:3.13-slim
ENV PYTHONUNBUFFERED=1
WORKDIR /app
COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt
COPY tick.py .
CMD ["python", "tick.py"]
EOF
cat > compose.yaml <<'EOF'
services:
redis:
image: redis:8-alpine
command: ["redis-server", "--save", "", "--appendonly", "no"]
healthcheck:
test: ["CMD", "redis-cli", "ping"]
interval: 3s
retries: 5
tick:
build: .
environment:
REDIS_URL: "redis://redis:6379/0"
INTERVAL: "3"
depends_on:
redis:
condition: service_healthy
stop_grace_period: 10s
deploy:
replicas: 3
EOF
docker compose up -d --build --scale tick=3 > /dev/null 2>&1
sleep 14
echo "═══ три реплики, задача должна выполняться один раз за интервал ═══"
docker compose logs tick --no-log-prefix 2>/dev/null | grep -c "выполняю" | \
xargs printf ' выполнений за ~14 c при интервале 3 c: %s\n'
docker compose logs tick --no-log-prefix 2>/dev/null | grep -c "пропускаю" | \
xargs printf ' пропусков (блокировка сработала): %s\n'
docker compose down -v > /dev/null 2>&1
cd /tmp && rm -rf /tmp/sched
Ожидаемый вывод:
═══ три реплики, задача должна выполняться один раз за интервал ═══
выполнений за ~14 c при интервале 3 c: 4
пропусков (блокировка сработала): 8
Четыре выполнения за четыре интервала — ровно то, что нужно. Восемь пропусков — две реплики из трёх каждый раз получали отказ по блокировке.
Без SET ... NX было бы 12 выполнений: каждая реплика выполняла бы задачу самостоятельно.
Обратите внимание на срок жизни ключа: он меньше интервала. Если бы блокировка жила дольше интервала, следующий запуск оказался бы заблокирован остатком предыдущей.
Практическое упражнение
Задание. Доработайте resources/examples/worker-redis/ так, чтобы задача не терялась при SIGKILL.
Требования:
- Забирать задачу через
BLMOVEв списокprocessing, а не черезBLPOP. - После успешной обработки удалять задачу из
processing. - Обработчик идемпотентен: повторная обработка той же задачи не даёт побочного эффекта дважды.
- Отдельная команда возвращает «зависшие» задачи из
processingв очередь. - Graceful shutdown сохранён:
SIGTERMдозавершает текущую задачу, код выхода0. - Доказать: после
SIGKILLв разгар обработки задача не потеряна и выполняется ровно один раз.
Подсказки
Подсказка 1
BLMOVE source destination LEFT RIGHT timeout — в Python: client.blmove(src, dst, timeout, "LEFT", "RIGHT"). Возвращает саму задачу или None по таймауту.
Подсказка 2
Удаление конкретного элемента из списка: LREM processing 1 <значение>. Значение должно совпадать побайтово — храните ту же строку, что получили.
Подсказка 3
Идемпотентность: SET result:{id} ... NX возвращает True только первому. Если False — задача уже обработана, побочный эффект повторять не нужно.
Подсказка 4
Для требования 6 нужен счётчик побочного эффекта в Redis: сравните его значение с числом задач.
Решение
Сначала выполните задание самостоятельно.
Показать решение
mkdir -p /tmp/wk-rel/worker && cd /tmp/wk-rel
cat > worker/__init__.py <<'PY'
"""Надёжный worker."""
PY
cat > worker/main.py <<'PY'
"""Worker с надёжной выборкой задач и идемпотентной обработкой.
Отличия от простого варианта на BLPOP:
* задача атомарно перемещается в processing и не исчезает при падении;
* повторная обработка безвредна благодаря ключу идемпотентности;
* зависшие задачи возвращаются в очередь отдельной командой.
"""
from __future__ import annotations
import json
import logging
import os
import signal
import sys
import time
from typing import Any
import redis
logging.basicConfig(
level=os.environ.get("LOG_LEVEL", "INFO"),
format="%(asctime)s %(levelname)s [worker] %(message)s",
stream=sys.stdout,
)
logger = logging.getLogger("worker")
QUEUE = os.environ.get("QUEUE_NAME", "tasks")
PROCESSING = f"{QUEUE}:processing"
EFFECT_COUNTER = f"{QUEUE}:effects"
POP_TIMEOUT = int(os.environ.get("POP_TIMEOUT", "2"))
DONE_TTL = int(os.environ.get("DONE_TTL", "3600"))
_shutdown = False
def request_shutdown(signum: int, _frame: object) -> None:
global _shutdown
logger.info("получен %s — завершаю после текущей задачи",
signal.Signals(signum).name)
_shutdown = True
def process(client: redis.Redis, task: dict[str, Any]) -> None:
"""Идемпотентная обработка (требование 3).
Побочный эффект — инкремент счётчика — выполняется только если
ключ идемпотентности удалось поставить. Повтор ничего не добавит.
"""
task_id = task.get("id", "?")
duration = float(task.get("duration", 0.5))
first_time = client.set(f"{QUEUE}:done:{task_id}", "1", nx=True, ex=DONE_TTL)
if not first_time:
logger.info("задача %s: уже обработана, пропускаю побочный эффект", task_id)
return
logger.info("задача %s: начата", task_id)
time.sleep(min(duration, 30.0))
client.incr(EFFECT_COUNTER) # побочный эффект
logger.info("задача %s: завершена", task_id)
def requeue(client: redis.Redis) -> int:
"""Требование 4: вернуть зависшие задачи из processing в очередь."""
moved = 0
while client.lmove(PROCESSING, QUEUE, "LEFT", "RIGHT"):
moved += 1
return moved
def wait_for_redis(client: redis.Redis) -> bool:
for attempt in range(1, 31):
try:
client.ping()
return True
except redis.ConnectionError:
logger.warning("Redis недоступен, попытка %s/30", attempt)
time.sleep(1)
return False
def main(argv: list[str] | None = None) -> int:
argv = sys.argv[1:] if argv is None else argv
client = redis.Redis.from_url(
os.environ.get("REDIS_URL", "redis://localhost:6379/0"),
decode_responses=True,
)
if not wait_for_redis(client):
logger.error("Redis не отвечает — выхожу")
return 1
if argv and argv[0] == "requeue":
moved = requeue(client)
logger.info("возвращено в очередь задач: %s", moved)
return 0
signal.signal(signal.SIGTERM, request_shutdown)
signal.signal(signal.SIGINT, request_shutdown)
logger.info("подключён к Redis, очередь %r, PID=%s", QUEUE, os.getpid())
processed = 0
while not _shutdown:
# Требование 1: атомарное перемещение вместо удаления.
# При падении здесь задача остаётся в PROCESSING.
raw = client.blmove(QUEUE, PROCESSING, POP_TIMEOUT, "LEFT", "RIGHT")
if raw is None:
continue
try:
task = json.loads(raw)
except json.JSONDecodeError:
logger.error("некорректный JSON, задача отброшена: %r", raw[:80])
client.lrem(PROCESSING, 1, raw)
continue
process(client, task)
# Требование 2: снять задачу с обработки только после успеха
client.lrem(PROCESSING, 1, raw)
processed += 1
logger.info("остановлен штатно, обработано задач: %s", processed)
return 0
if __name__ == "__main__":
sys.exit(main())
PY
cat > producer.py <<'PY'
"""Наполнение очереди задачами."""
from __future__ import annotations
import json
import os
import sys
import redis
def main() -> int:
count = int(sys.argv[1]) if len(sys.argv) > 1 else 5
duration = float(sys.argv[2]) if len(sys.argv) > 2 else 1.0
client = redis.Redis.from_url(
os.environ.get("REDIS_URL", "redis://localhost:6379/0"),
decode_responses=True,
)
queue = os.environ.get("QUEUE_NAME", "tasks")
for i in range(1, count + 1):
client.rpush(queue, json.dumps({"id": i, "duration": duration}))
print(f"поставлено задач: {count} (по {duration} c каждая)")
return 0
if __name__ == "__main__":
sys.exit(main())
PY
cat > requirements.txt <<'EOF'
redis==8.1.0
EOF
cat > .dockerignore <<'EOF'
Dockerfile
.dockerignore
__pycache__
*.py[cod]
EOF
cat > Dockerfile <<'EOF'
# syntax=docker/dockerfile:1
FROM python:3.13-slim
ENV PYTHONUNBUFFERED=1 \
PYTHONDONTWRITEBYTECODE=1
RUN useradd --create-home --uid 10001 appuser
WORKDIR /app
COPY requirements.txt .
RUN --mount=type=cache,target=/root/.cache/pip pip install -r requirements.txt
COPY --chown=appuser:appuser worker/ ./worker/
COPY --chown=appuser:appuser producer.py .
USER 10001:10001
CMD ["python", "-m", "worker.main"]
EOF
cat > compose.yaml <<'EOF'
services:
redis:
image: redis:8-alpine
command: ["redis-server", "--save", "", "--appendonly", "no"]
healthcheck:
test: ["CMD", "redis-cli", "ping"]
interval: 3s
retries: 5
worker:
build: .
environment:
REDIS_URL: "redis://redis:6379/0"
QUEUE_NAME: "tasks"
depends_on:
redis:
condition: service_healthy
# Больше самой длинной задачи
stop_grace_period: 40s
EOF
# ── Проверка требований ──
docker compose up -d --build > /dev/null 2>&1
sleep 10
echo "═══ Требование 5: graceful shutdown ═══"
docker compose exec -T worker python producer.py 3 3 > /dev/null
sleep 4
s="$(date +%s.%N)"; docker compose stop worker > /dev/null; e="$(date +%s.%N)"
printf ' время остановки: %.1f c, код: %s\n' \
"$(awk -v a="$s" -v b="$e" 'BEGIN{print b-a}')" \
"$(docker compose ps -a --format '{{.ExitCode}}' worker)"
docker compose logs worker --no-log-prefix 2>/dev/null | tail -2 | sed 's/^/ /'
echo
echo "═══ Требование 6: задача не теряется при SIGKILL ═══"
docker compose exec -T redis redis-cli flushall > /dev/null
docker compose up -d worker > /dev/null 2>&1
sleep 5
docker compose exec -T worker python producer.py 3 10 > /dev/null
sleep 4
printf ' до SIGKILL: очередь=%s processing=%s\n' \
"$(docker compose exec -T redis redis-cli llen tasks)" \
"$(docker compose exec -T redis redis-cli llen tasks:processing)"
docker compose kill -s SIGKILL worker > /dev/null 2>&1
sleep 1
printf ' после SIGKILL: очередь=%s processing=%s ← задача цела\n' \
"$(docker compose exec -T redis redis-cli llen tasks)" \
"$(docker compose exec -T redis redis-cli llen tasks:processing)"
echo
echo "═══ Требование 4: возврат зависших задач ═══"
docker compose run --rm -T worker python -m worker.main requeue 2>&1 | tail -1 | sed 's/^/ /'
printf ' после requeue: очередь=%s processing=%s\n' \
"$(docker compose exec -T redis redis-cli llen tasks)" \
"$(docker compose exec -T redis redis-cli llen tasks:processing)"
echo
echo "═══ Требование 3: идемпотентность ═══"
docker compose up -d worker > /dev/null 2>&1
sleep 40
printf ' задач поставлено: 3\n'
printf ' побочных эффектов: %s\n' \
"$(docker compose exec -T redis redis-cli get tasks:effects)"
printf ' осталось в processing: %s\n' \
"$(docker compose exec -T redis redis-cli llen tasks:processing)"
docker compose logs worker --no-log-prefix 2>/dev/null | grep -c "уже обработана" | \
xargs printf ' повторов, распознанных как дубли: %s\n'
docker compose down -v > /dev/null 2>&1
cd /tmp && rm -rf /tmp/wk-rel
Ожидаемый вывод:
═══ Требование 5: graceful shutdown ═══
время остановки: 2.1 c, код: 0
задача 2: завершена
остановлен штатно, обработано задач: 2
═══ Требование 6: задача не теряется при SIGKILL ═══
до SIGKILL: очередь=2 processing=1
после SIGKILL: очередь=2 processing=1 ← задача цела
═══ Требование 4: возврат зависших задач ═══
возвращено в очередь задач: 1
после requeue: очередь=3 processing=0
═══ Требование 3: идемпотентность ═══
задач поставлено: 3
побочных эффектов: 3
осталось в processing: 0
повторов, распознанных как дубли: 1
Все шесть требований выполнены. Разберём главные строки.
processing=1 после SIGKILL. Это и есть решение исходной проблемы: в варианте с BLPOP задача исчезала бесследно, здесь она видна и восстановима.
побочных эффектов: 3 при трёх задачах и одном повторе. Задача, прерванная SIGKILL, вернулась в очередь и была обработана заново — то есть выполнена дважды. Но счётчик показывает 3, а не 4: ключ идемпотентности не дал повторить побочный эффект. Строка «повторов, распознанных как дубли: 1» подтверждает, что дубль действительно случился и был распознан.
Это ровно то, ради чего нужна идемпотентность: система даёт at-least-once, обработчик делает повтор безвредным.
Три решения, определяющие качество.
LREM вызывается после process, а не до. Порядок здесь — весь смысл: задача покидает processing только после успешного завершения. Обратный порядок вернул бы поведение BLPOP.
Ключ идемпотентности ставится в начале обработки, а не в конце. Так дубль отсекается до побочного эффекта. Обратите внимание на цену этого решения: если worker умрёт между установкой ключа и инкрементом, задача при повторе будет пропущена — эффект потеряется. Это осознанный размен: защита от дубля важнее, потому что дубль обычно вреднее пропуска (двойное списание против неотправленного письма). Где пропуск дороже — ключ ставят после эффекта, принимая риск дубля.
requeue — отдельная команда того же образа, а не отдельный сервис. Один образ, одна кодовая база, запуск через docker compose run. Для production её вызывают по расписанию с блокировкой из раздела о периодических задачах.
Чего решение не делает. Нет ограничения на число повторов: задача, падающая всегда, будет возвращаться в очередь бесконечно. В production нужен счётчик попыток и очередь «мёртвых» задач (dead letter queue). Нет и определения «зависшей» задачи по времени — requeue возвращает всё содержимое processing, поэтому запускать её на работающих worker'ах нельзя. Обе задачи решены в готовых библиотеках (Celery, RQ, Dramatiq) — что и есть основная причина их использовать.
Проверка результата
cd resources/examples/worker-redis
docker compose up -d --build
sleep 12
docker compose exec -T worker python producer.py 3 2
sleep 3
docker compose stop worker
docker compose ps -a --format 'table {{.Service}}\t{{.ExitCode}}'
docker compose logs worker --no-log-prefix | tail -3
docker compose down -v
Ожидается код выхода 0, строка «остановлен штатно» и завершённая, а не оборванная задача.
Типичные ошибки
| Ошибка | Причина | Исправление |
|---|---|---|
workers = 2 * os.cpu_count() + 1 | Формула из документации Gunicorn | В container даёт CPU хоста; задать WEB_CONCURRENCY явно |
os.cpu_count() для расчёта при --cpus | Кажется очевидным | Читать /sys/fs/cgroup/cpu.max |
--threads для CPU-bound нагрузки | Потоки кажутся дешевле процессов | GIL не даст выигрыша; нужны процессы |
| Лимит памяти без учёта числа worker'ов | Считали по одному процессу | Умножить RSS на число worker'ов, добавить 30 % |
Соединение с базой при импорте плюс preload_app | Кажется оптимизацией | Дескриптор наследуется всеми worker'ами; создавать после fork |
| Блокирующее ожидание задачи без таймаута | Проще писать | Цикл не проверит флаг завершения; висит до SIGKILL |
| Обработка задачи внутри обработчика сигнала | Кажется удобным | Обработчик прерывает поток в произвольной точке; только флаг |
stop_grace_period по умолчанию для worker'а | Не задумывались | 10 секунд мало для длинной задачи; задать явно |
| Неидемпотентный обработчик | Расчёт на exactly-once | Гарантия — at-least-once; добавить ключ идемпотентности |
BLPOP в надёжной очереди | Простейший вариант | Задача теряется при падении; BLMOVE в processing |
cron внутри container с приложением | Привычка с виртуальных машин | Планировщик оркестратора или отдельный container |
| Планировщик в нескольких репликах без блокировки | Не учли масштабирование | SET key NX EX перед выполнением |
Контрольные вопросы
На понимание:
- Почему
os.cpu_count()не видит--cpus, но видит--cpuset-cpus? - Почему потоки не ускоряют CPU-bound нагрузку в Python?
- Что даёт
preload_appи почему экономия памяти меньше ожидаемой? - Почему очереди дают at-least-once, а не exactly-once?
- Чем
BLMOVEнадёжнееBLPOPи какой ценой?
На применение:
- Как рассчитать
--memoryдля 4 worker'ов Flask? - Как сделать обработчик задач идемпотентным?
- Как организовать ежедневную задачу при трёх репликах сервиса?
На диагностику:
- Container с 8 worker'ами завершается с кодом
137. Порядок диагностики? - Worker игнорирует
docker stopи убивается через 10 секунд. Две вероятные причины?
Краткое резюме
- Формула
2 × CPU + 1в container даёт число по CPU хоста и приводит к OOM. --cpusне виден ниos.cpu_count(), ниsched_getaffinity; правду знает/sys/fs/cgroup/cpu.max.--cpuset-cpus, в отличие от--cpus, меняет маску affinity и потому виден.- В production число worker'ов задаётся явно через
WEB_CONCURRENCY, а не вычисляется. - Потоки помогают только I/O-bound нагрузке; CPU-bound требует процессов.
- Память растёт почти линейно с числом worker'ов; лимит считается как RSS × N × 1.3.
preload_appускоряет старт, но требует создавать соединения послеfork.- Background worker обязан ждать задачу с таймаутом, иначе не увидит флаг завершения.
- Обработчик сигнала только выставляет флаг; работу доводит основной цикл.
- Очереди дают at-least-once — обработчик обязан быть идемпотентным.
BLPOPтеряет задачу при падении worker'а;BLMOVEвprocessing— нет.cronвнутри container — антипаттерн; при нескольких репликах нужна распределённая блокировка.
Официальные источники
| Источник | Ссылка | Что подтверждает |
|---|---|---|
| Gunicorn: design | https://docs.gunicorn.org/en/stable/design.html | Pre-fork модель, формула 2 × cores + 1, классы worker'ов |
| Gunicorn: settings | https://docs.gunicorn.org/en/stable/settings.html | preload_app, threads, WEB_CONCURRENCY, хуки when_ready и post_fork |
| Uvicorn: deployment | https://www.uvicorn.org/deployment/ | Управление worker'ами, --workers |
| Docker: resource constraints | https://docs.docker.com/engine/containers/resource_constraints/ | --cpus, --cpuset-cpus, --memory, --pids-limit |
| cgroup v2 | https://docs.kernel.org/admin-guide/cgroup-v2.html | Формат cpu.max, квота и период |
Python: os.cpu_count и sched_getaffinity | https://docs.python.org/3/library/os.html#os.cpu_count | Что именно возвращает каждый вызов |
Python: gc.freeze | https://docs.python.org/3/library/gc.html#gc.freeze | Уменьшение copy-on-write после fork |
Redis: BLMOVE | https://redis.io/docs/latest/commands/blmove/ | Атомарное перемещение, замена BRPOPLPUSH |
Redis: SET с NX и EX | https://redis.io/docs/latest/commands/set/ | Распределённая блокировка |
Compose: stop_grace_period | https://docs.docker.com/reference/compose-file/services/ | Время до SIGKILL |
| Celery: worker | https://docs.celeryq.dev/en/stable/userguide/workers.html | at-least-once, повторы, graceful shutdown |
Навигация
← Предыдущий материал
Вернуться к разделу
Следующий материал → Тестирование в container
Главное оглавление