Главная/Python внутри Container/Урок

6.11. Worker processes

Цели

После этого материала вы сможете:

  • рассчитать число worker-процессов исходя из лимита container, а не из числа CPU хоста;
  • объяснить, почему os.cpu_count() даёт неверный ответ при --cpus;
  • выбрать между процессами и потоками для вашего типа нагрузки;
  • оценить расход памяти на worker и понять, что даёт preload;
  • написать background worker, который корректно завершается по SIGTERM и не теряет задачи;
  • объяснить, почему обработчик задач обязан быть идемпотентным;
  • организовать периодические задачи без cron внутри container.

Предварительные знания

Рабочий пример — 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) + 1 as the number of workers to start off with.

Формула разумна на выделенном сервере. В container она даёт неверный результат, потому что «num_cores» обычно получают так:

python
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
Рост latency33 процесса делят квоту 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:

python
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 предпочтительнее противоположный подход — задать число явно переменной окружения:

yaml
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

Каждый процесс — отдельный интерпретатор со своей копией импортированных модулей. Базовый расход:

КомпонентПорядок
Интерпретатор Python10–15 MB
Flask плюс зависимости+20–30 MB
FastAPI плюс Pydantic+40–60 MB
Модели ML, кэши, данныезависит от приложения

Отсюда правило подбора лимита памяти:

text
--memory ≥ (RSS одного worker'а × число worker'ов) × 1.3

Запас в 30 % покрывает фрагментацию аллокатора и пики при обработке запроса (урок 6.13).

Что даёт preload

Без preload каждый worker импортирует приложение сам после fork. С preload master импортирует приложение один раз, а fork наследует уже загруженную память.

python
# gunicorn.conf.py
preload_app = True

Ожидаемая выгода — copy-on-write: после fork страницы физически общие, копируются только при записи. Значит, код и константы приложения существуют в одном экземпляре.

Но в CPython эта выгода частично теряется. Причина — счётчики ссылок: они хранятся в заголовке каждого объекта. Любое обращение к объекту меняет ob_refcnt, то есть пишет в его страницу, и ядро копирует страницу целиком. Сборщик мусора усугубляет: он обходит объекты и трогает их все.

Смягчение — gc.freeze(), доступный с Python 3.7:

python
# 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 workerBackground worker
Входящие соединенияЕстьНет
ПортПубликуетсяНе нужен
HEALTHCHECKHTTP-запросНет естественного endpoint
Единица работыЗапрос, миллисекундыЗадача, секунды или минуты
Grace period10 секунд обычно хватаетЧасто нужно больше
Потеря при остановкеОдин запросЦелая задача

Ключевое требование — цикл должен проверять флаг завершения:

python
_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)

Два обязательных элемента:

  1. Обработчик только выставляет флаг. Обработчик сигнала прерывает основной поток в произвольной точке; сложная работа в нём небезопасна (урок 6.5).
  2. Ожидание задачи с таймаутом. Блокирующее ожидание без таймаута не даст циклу проверить флаг: worker будет висеть до SIGKILL.

Grace period задаётся отдельно, так как задача длиннее запроса:

yaml
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 вместо insertINSERT ... ON CONFLICT DO NOTHING
Проверка состоянияПеред отправкой письма проверить, не отправлено ли

Надёжность очереди на Redis

Простейший способ забрать задачу — BLPOP: он удаляет элемент из списка. Если worker будет убит после BLPOP, но до завершения обработки, задача исчезнет.

Надёжный вариант — атомарно переместить задачу в список «в работе»:

text
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Простые интервальные задачи

Общее требование ко всем: защита от одновременного выполнения. При нескольких репликах планировщик запустится в каждой. Решение — распределённая блокировка:

text
SET cron:daily-report locked NX EX 3600

Только один экземпляр получит OK и выполнит задачу.


Внутренний механизм

Что происходит при fork

  1. Ядро создаёт копию процесса; таблицы страниц указывают на те же физические страницы, помеченные read-only.
  2. При записи в такую страницу возникает page fault.
  3. Ядро копирует страницу и снимает пометку.

Отсюда «copy-on-write»: копируется не всё, а только изменённое. В CPython изменяется больше, чем кажется, — из-за refcount.

Наследуются: открытые файловые дескрипторы, отображённая память, обработчики сигналов. Не наследуются: потоки (в дочернем остаётся только вызвавший fork).

Последнее — причина частой ошибки: пул соединений с фоновым потоком, созданный до fork, в worker'е окажется без своего потока и перестанет работать.

Почему --cpus невидим для процесса

--cpus 1.5 записывает в cgroup:

text
/sys/fs/cgroup/cpu.max → "150000 100000"

Планировщик ядра учитывает эту квоту, но интерфейс sched_getaffinity возвращает маску разрешённых CPU, а не квоту. Маска не менялась. Программе доступна только сама cgroup-файловая система — что и делает функция available_cpus() выше.


Команды и примеры

Три ответа на вопрос «сколько CPU»

bash
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-ядерной машине:

text
═══ без лимита ═══
  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

bash
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

Ожидаемый вывод (порядок величины; конкретные числа зависят от машины):

text
  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'ов

bash
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

Ожидаемый вывод:

text
═══ 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:

bash
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

Ожидаемый вывод:

text
процессов python: 4

Четыре процесса — менеджер и три worker'а: переменная учтена. Если у вас получилось 1, ваша версия Uvicorn её не читает — тогда задавайте --workers явно. Это ровно тот случай, когда стоит проверить командой, а не поверить тексту.

Потоки против процессов на CPU-bound нагрузке

bash
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"

Ожидаемый вывод (относительные значения важнее абсолютных):

text
═══ 8 параллельных CPU-bound запросов, --cpus 2 ═══
  1 процесс                    4.10 c
  1 процесс, 4 потока          4.18 c
  2 процесса                   2.15 c
  4 процесса                   2.20 c

Три вывода из этого измерения:

  1. Потоки не помогли — 4.18 против 4.10 секунды. GIL не даёт двум потокам исполнять байткод одновременно; появились лишь накладные расходы на переключение.
  2. Процессы помогли вдвое — у каждого свой GIL, и оба использовали квоту 2 CPU.
  3. Четыре процесса не дали ничего сверх двух — квота исчерпана на двух. Дальнейшее увеличение только тратит память.

Последнее — практическое доказательство того, что число worker'ов должно исходить из лимита container.

bash
docker rmi -f wk wkcpu > /dev/null 2>&1; rm -rf /tmp/wk /tmp/cpus.py

Background worker: graceful shutdown

bash
cd resources/examples/worker-redis
docker compose up -d --build
sleep 12
docker compose ps --format 'table {{.Service}}\t{{.Status}}'

Ожидаемый вывод:

text
SERVICE   STATUS
redis     Up 12 seconds (healthy)
worker    Up 6 seconds

Поставим задачи и остановим worker в момент обработки:

bash
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

Ожидаемый вывод:

text
поставлено задач: 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.

bash
printf 'осталось в очереди: %s\n' \
    "$(docker compose exec -T redis redis-cli llen tasks)"

Ожидаемый вывод:

text
осталось в очереди: 3

Почему BLPOP теряет задачу

Проверим, что происходит при жёстком убийстве:

bash
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 "  задача, которая обрабатывалась, нигде не осталась — потеряна"

Ожидаемый вывод:

text
═══ до SIGKILL ═══
  в очереди: 2
═══ после SIGKILL ═══
  в очереди: 2
  код выхода: 137
  задача, которая обрабатывалась, нигде не осталась — потеряна

Задач было три: одну worker забрал (BLPOP удалил её из списка), две остались. После SIGKILL в очереди по-прежнему две — забранная задача исчезла вместе с процессом.

Это не ошибка примера, а неизбежное свойство BLPOP. Устранение — в упражнении.

bash
docker compose down -v

Периодическая задача с блокировкой

bash
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

Ожидаемый вывод:

text
═══ три реплики, задача должна выполняться один раз за интервал ═══
  выполнений за ~14 c при интервале 3 c: 4
  пропусков (блокировка сработала):        8

Четыре выполнения за четыре интервала — ровно то, что нужно. Восемь пропусков — две реплики из трёх каждый раз получали отказ по блокировке.

Без SET ... NX было бы 12 выполнений: каждая реплика выполняла бы задачу самостоятельно.

Обратите внимание на срок жизни ключа: он меньше интервала. Если бы блокировка жила дольше интервала, следующий запуск оказался бы заблокирован остатком предыдущей.


Практическое упражнение

Задание. Доработайте resources/examples/worker-redis/ так, чтобы задача не терялась при SIGKILL.

Требования:

  1. Забирать задачу через BLMOVE в список processing, а не через BLPOP.
  2. После успешной обработки удалять задачу из processing.
  3. Обработчик идемпотентен: повторная обработка той же задачи не даёт побочного эффекта дважды.
  4. Отдельная команда возвращает «зависшие» задачи из processing в очередь.
  5. Graceful shutdown сохранён: SIGTERM дозавершает текущую задачу, код выхода 0.
  6. Доказать: после 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: сравните его значение с числом задач.

Решение

Сначала выполните задание самостоятельно.

Показать решение
bash
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

Ожидаемый вывод:

text
═══ Требование 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) — что и есть основная причина их использовать.

Проверка результата

bash
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 перед выполнением

Контрольные вопросы

На понимание:

  1. Почему os.cpu_count() не видит --cpus, но видит --cpuset-cpus?
  2. Почему потоки не ускоряют CPU-bound нагрузку в Python?
  3. Что даёт preload_app и почему экономия памяти меньше ожидаемой?
  4. Почему очереди дают at-least-once, а не exactly-once?
  5. Чем BLMOVE надёжнее BLPOP и какой ценой?

На применение:

  1. Как рассчитать --memory для 4 worker'ов Flask?
  2. Как сделать обработчик задач идемпотентным?
  3. Как организовать ежедневную задачу при трёх репликах сервиса?

На диагностику:

  1. Container с 8 worker'ами завершается с кодом 137. Порядок диагностики?
  2. Worker игнорирует docker stop и убивается через 10 секунд. Две вероятные причины?

Краткое резюме

  1. Формула 2 × CPU + 1 в container даёт число по CPU хоста и приводит к OOM.
  2. --cpus не виден ни os.cpu_count(), ни sched_getaffinity; правду знает /sys/fs/cgroup/cpu.max.
  3. --cpuset-cpus, в отличие от --cpus, меняет маску affinity и потому виден.
  4. В production число worker'ов задаётся явно через WEB_CONCURRENCY, а не вычисляется.
  5. Потоки помогают только I/O-bound нагрузке; CPU-bound требует процессов.
  6. Память растёт почти линейно с числом worker'ов; лимит считается как RSS × N × 1.3.
  7. preload_app ускоряет старт, но требует создавать соединения после fork.
  8. Background worker обязан ждать задачу с таймаутом, иначе не увидит флаг завершения.
  9. Обработчик сигнала только выставляет флаг; работу доводит основной цикл.
  10. Очереди дают at-least-once — обработчик обязан быть идемпотентным.
  11. BLPOP теряет задачу при падении worker'а; BLMOVE в processing — нет.
  12. cron внутри container — антипаттерн; при нескольких репликах нужна распределённая блокировка.

Официальные источники

ИсточникСсылкаЧто подтверждает
Gunicorn: designhttps://docs.gunicorn.org/en/stable/design.htmlPre-fork модель, формула 2 × cores + 1, классы worker'ов
Gunicorn: settingshttps://docs.gunicorn.org/en/stable/settings.htmlpreload_app, threads, WEB_CONCURRENCY, хуки when_ready и post_fork
Uvicorn: deploymenthttps://www.uvicorn.org/deployment/Управление worker'ами, --workers
Docker: resource constraintshttps://docs.docker.com/engine/containers/resource_constraints/--cpus, --cpuset-cpus, --memory, --pids-limit
cgroup v2https://docs.kernel.org/admin-guide/cgroup-v2.htmlФормат cpu.max, квота и период
Python: os.cpu_count и sched_getaffinityhttps://docs.python.org/3/library/os.html#os.cpu_countЧто именно возвращает каждый вызов
Python: gc.freezehttps://docs.python.org/3/library/gc.html#gc.freezeУменьшение copy-on-write после fork
Redis: BLMOVEhttps://redis.io/docs/latest/commands/blmove/Атомарное перемещение, замена BRPOPLPUSH
Redis: SET с NX и EXhttps://redis.io/docs/latest/commands/set/Распределённая блокировка
Compose: stop_grace_periodhttps://docs.docker.com/reference/compose-file/services/Время до SIGKILL
Celery: workerhttps://docs.celeryq.dev/en/stable/userguide/workers.htmlat-least-once, повторы, graceful shutdown

Навигация

← Предыдущий материал
Вернуться к разделу
Следующий материал → Тестирование в container
Главное оглавление

Markdown на GitHub ↗