DevOps & Scaling

Автомасштабирование пула воркеров для решения CAPTCHA

Сколько воркеров держать в пуле, если нагрузка за час может вырасти в десять раз, а потом так же резко упасть? Правильный ответ — не держать фиксированное число вообще, а считать глубину очереди, загрузку и остаток баланса и подстраивать пул под них автоматически. Статический пул почти всегда либо переплачивает в затишье, либо упирается в потолок на пике: очередь растёт, P95 по времени решения ползёт вверх, а часть задач не укладывается в таймаут вызывающего сервиса. Ниже — три рабочие схемы на Python (потоки, процессы, контроль баланса) и объяснение, когда вместо самодельного скейлера логичнее взять Kubernetes HPA или KEDA.


По каким сигналам масштабировать пул

Сигналы, которые двигают пул сразу

  • Глубина очереди. Расширяйте пул при > 20 задачах в ожидании, сокращайте при < 5. Реагирует быстрее всех остальных сигналов.
  • Загрузка воркеров. Рост — когда заняты > 80 %, сжатие — когда заняты < 20 %. Второй по скорости реакции сигнал.

Сигналы-индикаторы, что с пулом что-то не так

  • Задержка решения (P95). P95 > 60 секунд говорит, что пула не хватает, даже если очередь ещё не выглядит критично; P95 < 20 секунд — можно сжиматься.
  • Доля ошибок. Рост доли ошибок выше 5 % обычно означает не нехватку воркеров, а протухшие сессии или устаревший код решения — нужен не столько новый воркер, сколько перезапуск текущих. Стабильные < 1 % — сигнал в порядке.
  • Баланс. Работает только в одну сторону: никогда не заставляет пул расти, а лишь останавливает расширение при балансе ниже $1.

Автомасштабирование на потоках

Если воркер только вызывает API CaptchaAI и ждёт ответ — это чистая I/O-нагрузка, и масштабировать в рамках одного процесса потоками дешевле и проще, чем поднимать отдельные процессы. Ниже — пул, который сам добавляет и убирает потоки раз в 10 секунд, ориентируясь на глубину очереди Redis и текущую занятость:

import os
import time
import threading
import requests
import json
import redis


class AutoScalingPool:
    """Dynamically scale CaptchaAI worker threads."""

    def __init__(self, api_key, redis_url="redis://localhost:6379"):
        self.api_key = api_key
        self.redis = redis.from_url(redis_url)
        self.base = "https://ocr.captchaai.com"
        self.queue_key = "captcha:tasks"
        self.results_key = "captcha:results"

        self.min_workers = 2
        self.max_workers = 20
        self.workers = []
        self.active_count = 0
        self.lock = threading.Lock()
        self.running = True

    def start(self):
        """Start the pool with minimum workers."""
        for _ in range(self.min_workers):
            self._add_worker()

        # Start scaler in background
        scaler = threading.Thread(target=self._scaling_loop, daemon=True)
        scaler.start()
        print(f"Pool started with {self.min_workers} workers")

    def _add_worker(self):
        """Add a worker thread."""
        if len(self.workers) >= self.max_workers:
            return
        t = threading.Thread(target=self._worker_loop, daemon=True)
        t.start()
        self.workers.append(t)

    def _remove_worker(self):
        """Signal one worker to stop (lazy removal)."""
        if len(self.workers) <= self.min_workers:
            return
        self.workers.pop()  # Thread will exit on next idle cycle

    def _worker_loop(self):
        """Worker loop: fetch and process tasks."""
        while self.running and threading.current_thread() in self.workers:
            result = self.redis.blpop(self.queue_key, timeout=10)
            if result is None:
                continue

            _, raw = result
            task = json.loads(raw)
            task_id = task["id"]

            with self.lock:
                self.active_count += 1

            try:
                token = self._solve(task["method"], task["params"])
                self.redis.hset(self.results_key, task_id, json.dumps({
                    "status": "success", "token": token,
                }))
            except Exception as e:
                self.redis.hset(self.results_key, task_id, json.dumps({
                    "status": "error", "error": str(e),
                }))
            finally:
                with self.lock:
                    self.active_count -= 1

    def _scaling_loop(self):
        """Periodically adjust worker count."""
        while self.running:
            time.sleep(10)

            queue_depth = self.redis.llen(self.queue_key)
            current = len(self.workers)
            utilization = (
                self.active_count / current * 100 if current > 0 else 0
            )

            # Scale up: queue growing and workers busy
            if queue_depth > 20 and utilization > 70:
                new_count = min(current + 2, self.max_workers)
                while len(self.workers) < new_count:
                    self._add_worker()
                print(f"Scaled up: {current} → {len(self.workers)} workers")

            # Scale down: queue empty and workers idle
            elif queue_depth < 5 and utilization < 20:
                target = max(current - 1, self.min_workers)
                while len(self.workers) > target:
                    self._remove_worker()
                if len(self.workers) < current:
                    print(f"Scaled down: {current} → {len(self.workers)} workers")

    def _solve(self, method, params, timeout=120):
        data = {"key": self.api_key, "method": method, "json": 1}
        data.update(params)

        resp = requests.post(
            f"{self.base}/in.php", data=data, timeout=30,
        )
        result = resp.json()

        if result.get("status") != 1:
            raise RuntimeError(result.get("request"))

        captcha_id = result["request"]
        start = time.time()

        while time.time() - start < timeout:
            time.sleep(5)
            resp = requests.get(f"{self.base}/res.php", params={
                "key": self.api_key,
                "action": "get",
                "id": captcha_id,
                "json": 1,
            }, timeout=15)
            data = resp.json()
            if data["request"] != "CAPCHA_NOT_READY":
                if data.get("status") == 1:
                    return data["request"]
                raise RuntimeError(data["request"])

        raise TimeoutError("Solve timeout")

    def stats(self):
        return {
            "workers": len(self.workers),
            "active": self.active_count,
            "queue": self.redis.llen(self.queue_key),
        }


# Usage
pool = AutoScalingPool(os.environ["CAPTCHAAI_KEY"])
pool.start()

# Monitor
while True:
    print(pool.stats())
    time.sleep(30)

Почему рост и сжатие несимметричны

Обратите внимание на асимметрию в _scaling_loop: пул расширяется сразу на два воркера, а сжимается только на одного за цикл — осознанная защита от «дребезга», когда пул мечется туда-обратно при малейшем колебании очереди. _remove_worker не убивает поток принудительно — он выводит его из списка активных, и поток завершается сам на следующей пустой итерации blpop, а задача, которую воркер уже начал решать, всегда доводится до конца.


Автомасштабирование на процессах

Как только в воркер добавляется что-то тяжелее сетевого вызова — предобработка изображений, локальный OCR, распаковка градиентных капч перед отправкой в API — GIL перестаёт быть безобидным, и потоки перестают давать реальный параллелизм на CPU-bound части работы. Здесь нужна изоляция на уровне процессов:

import multiprocessing
import time
import redis
import os


class ProcessScaler:
    """Scale worker processes based on queue depth."""

    def __init__(self, worker_fn, redis_url="redis://localhost:6379"):
        self.worker_fn = worker_fn
        self.redis = redis.from_url(redis_url)
        self.processes = []
        self.min_workers = 2
        self.max_workers = 16

    def run(self, check_interval=15):
        """Run the scaler loop."""
        # Start minimum workers
        for _ in range(self.min_workers):
            self._spawn()

        while True:
            time.sleep(check_interval)
            self._cleanup_dead()

            queue_depth = self.redis.llen("captcha:tasks")
            current = len(self.processes)

            # Scale up
            if queue_depth > current * 5 and current < self.max_workers:
                to_add = min(
                    max(1, queue_depth // 10),
                    self.max_workers - current,
                )
                for _ in range(to_add):
                    self._spawn()
                print(f"Scaled up to {len(self.processes)} workers")

            # Scale down
            elif queue_depth < 3 and current > self.min_workers:
                to_remove = min(2, current - self.min_workers)
                for _ in range(to_remove):
                    p = self.processes.pop()
                    p.terminate()
                print(f"Scaled down to {len(self.processes)} workers")

    def _spawn(self):
        p = multiprocessing.Process(target=self.worker_fn)
        p.start()
        self.processes.append(p)

    def _cleanup_dead(self):
        self.processes = [p for p in self.processes if p.is_alive()]
        # Ensure minimum
        while len(self.processes) < self.min_workers:
            self._spawn()

Почему порог роста относительный, а не абсолютный

Здесь порог роста считается относительно текущего размера пула (queue_depth > current * 5), а не абсолютным числом задач — это не даёт маленькому пулу переоценивать нагрузку и сразу прыгать на максимум. _cleanup_dead() стоит вызывать в начале каждого цикла: процессы, в отличие от потоков, могут упасть по OOM или необработанному исключению интерпретатора и просто исчезнуть из multiprocessing.Process.is_alive(), а пул должен это заметить и восстановить минимум сам.


Защита баланса от быстрого расхода

У автомасштабирования есть очевидный побочный эффект: чем больше воркеров, тем быстрее тратится баланс аккаунта. Без отдельной проверки скейлер, который реагирует только на очередь, будет наращивать пул до упора, даже если денег хватит на десять минут работы. Проверку баланса стоит выносить в отдельную функцию и дергать её перед каждым решением о расширении:

def check_balance(api_key, min_balance=2.0):
    """Check if balance is sufficient for scaling."""
    resp = requests.get("https://ocr.captchaai.com/res.php", params={
        "key": api_key,
        "action": "getbalance",
        "json": 1,
    }, timeout=15)
    balance = float(resp.json()["request"])

    if balance < min_balance:
        print(f"Balance ${balance:.2f} below ${min_balance} — halting scale-up")
        return False
    return True

Встраивается это одной проверкой прямо в цикл масштабирования:

# In _scaling_loop:
if queue_depth > 20 and utilization > 70:
    if check_balance(self.api_key, min_balance=2.0):
        # Scale up
        ...
    else:
        print("Scaling paused — low balance")

Как связать max_workers с тарифом CaptchaAI

Порог min_balance стоит привязывать не к произвольному числу, а к тарифу: биллинг у CaptchaAI идёт по одновременным потокам, а не по решённой задаче, и max_workers в коде выше фактически равен пиковому числу потоков, которое нужно от тарифа. Для честных 20 параллельных потоков нужен минимум ADVANCE ($90/мес, 50 потоков) с запасом, под пул на 100+ процессов — PREMIUM ($170/мес, 100 потоков) или CORPORATE ($240/мес, 150 потоков). Сверить max_workers с числом потоков в тарифе заранее дешевле, чем ловить нехватку потоков на проде.

Локальный сценарий: биллинг в валюте клиента

Для команд, которые выставляют счета клиентам в рублях, тенге или гривне, а платят провайдеру в долларах, фиксированная цена в USD за поток удобна ещё и тем, что стоимость масштабирования предсказуема заранее — не нужно пересчитывать курс на каждый всплеск нагрузки.


Типичные проблемы и что с ними делать

Пул растёт или сжимается неправильно

  1. Пул только растёт и не сжимается. Обычно значит, что очередь никогда не пустеет — проверьте, действительно ли воркеры обрабатывают задачи, а не падают молча.
  2. Сжатие пула слишком резкое. Порог сжатия занижен — увеличьте задержку перед сжатием до 30 секунд и более.

Проблемы с ресурсами и балансом

  1. Остаются процессы-зомби. Завершённые процессы не вычищаются — регулярно вызывайте _cleanup_dead().
  2. Баланс тает быстрее ожидаемого. Слишком много воркеров работает одновременно — добавьте проверку баланса в логику масштабирования и сверьте max_workers с тарифом.

При логировании метрик пула (id задач, IP источника трафика) собирайте только данные, которые вы вправе обрабатывать по 152-ФЗ «О персональных данных» или GDPR при трафике из ЕС — это обычная гигиена логов, а не особенность CaptchaAI.


Что выбрать: свой скейлер, Kubernetes HPA или KEDA

Четыре рабочих варианта, от самого лёгкого до самого тяжёлого:

  • Пул потоков — для I/O-bound нагрузки (только вызовы API). Задержка реакции низкая, сложность реализации низкая.
  • Пул процессов — для предобработки с нагрузкой на CPU. Задержка реакции средняя, сложность средняя.
  • Kubernetes HPA — для облачных развёртываний, которые уже живут в кластере. Задержка реакции выше, чем у самописных скейлеров, сложность внедрения высокая.
  • KEDA — для масштабирования по внешним метрикам (длина очереди Redis, Kafka). Задержка средняя, сложность средняя, но не нужен отдельный контроллер.

Если сервис уже в Kubernetes, логику скейлера часто не нужно писать с нуля — HPA масштабирует поды по кастомным метрикам, а KEDA даёт готовые адаптеры под очереди без лишнего кода. Самописный скейлер оправдан, когда воркеры живут вне кластера — например, на паре VPS в европейском регионе или в Казахстане с низким RTT до ocr.captchaai.com — и поднимать ради них Kubernetes избыточно.


Часто задаваемые вопросы

Сколько воркеров держать на одну очередь задач?

Ориентир — один воркер на 5–10 задач в очереди. Один воркер обрабатывает примерно 3–6 CAPTCHA в минуту в зависимости от типа: простые image-капчи решаются быстрее, GeeTest и Cloudflare Turnstile — медленнее.

Что выбрать — потоки или процессы?

Потоки — если воркер только дергает API CaptchaAI и ждёт ответ (чистый I/O). Процессы — если помимо решения воркер ещё делает предобработку изображений или другие CPU-тяжёлые операции.

Какой тариф CaptchaAI взять под автомасштабируемый пул?

Смотрите на пиковое значение max_workers, а не на среднее: тариф должен покрывать потоками именно пиковую одновременную нагрузку. Для пула до 50 параллельных воркеров хватает ADVANCE ($90/мес, 50 потоков); для более крупных всплесков — PREMIUM ($170/мес, 100 потоков) или CORPORATE ($240/мес, 150 потоков).

Что делать, если баланс заканчивается быстрее, чем ожидалось?

Это почти всегда значит, что пул расширяется без проверки баланса или что min_balance выставлен слишком низко для скорости расхода. Добавьте check_balance() в каждое решение о расширении и поднимите порог настолько, чтобы у вас было хотя бы 5–10 минут на реакцию до полного исчерпания средств.


Связанные руководства


Масштабируйте пул с умом — получите API-ключ CaptchaAI сегодня.

Комментарии для этой статьи отключены.