DevOps & Scaling

Обмен сообщениями NATS + CaptchaAI: упрощенное распределение задач CAPTCHA

Если нужно раскидать тысячи задач CAPTCHA по нескольким воркерам, но поднимать кластер Kafka с ZooKeeper ради этого не хочется — берите NATS. Один бинарник, задержка меньше миллисекунды, память в районе 20 МБ, а группы очередей (queue groups) сами балансируют нагрузку между обработчиками. Для очереди CAPTCHA-задач, где важна скорость, а не гарантии доставки Kafka-уровня, это ровно то, что нужно.

Архитектура: от скрапера до результата

[Scrapers] → Publish → [NATS: captcha.tasks]
                              ↓
                    Queue Group: captcha-workers
                    ├── Worker 1 (solve via CaptchaAI)
                    ├── Worker 2
                    └── Worker 3
                              ↓
                    Publish → [NATS: captcha.results]
                              ↓
                    [Result Subscribers]

Группы очередей NATS автоматически распределяют сообщения между воркерами — каждая задача достаётся ровно одному обработчику, даже если параллельно запущено несколько процессов-подписчиков на один и тот же subject.

Пример из практики: команда, распределённая между Минском и Алматы, собирает данные и прогоняет капчи через воркеров на разных серверах — NATS достаточно поднять один раз в общей сети, а очередь captcha-workers сама распределяет нагрузку между регионами. На плане ADVANCE ($90/мес, 50 потоков) такая связка спокойно тянет несколько сотен задач в минуту, пока воркеры не упираются в собственный CPU, а не в лимит CaptchaAI.

Когда NATS выигрывает у Kafka и RabbitMQ

Прежде чем встраивать очередь в продакшен, стоит понимать, чем NATS отличается от более тяжёлых альтернатив:

Характеристика NATS Kafka RabbitMQ
Задержка < 1 мс 5–10 мс 1–5 мс
Сложность развёртывания Один бинарник Кластер + ZooKeeper Средняя
Потребление памяти ~20 МБ ~1 ГБ+ ~200 МБ
Персистентность Опционально (JetStream) Встроена Встроена
Подходит для Недолговечных задач с низкой задержкой Надёжного стриминга Сложной маршрутизации

Задачи CAPTCHA живут секунды: если сообщение потерялось, вы просто отправляете задачу заново. Из этого следуют два практических вывода:

  • Терять сообщение некритично — переотправка стоит дешевле, чем поддержка Kafka-кластера ради гарантии доставки, которая здесь не нужна.
  • Скорость и простота эксплуатации важнее персистентности — один бинарник NATS поднимается за секунды и не требует отдельной команды для поддержки брокера.

Что понадобится: установка NATS и клиентов

Перед первой интеграцией нужны только сервер NATS и клиентская библиотека под ваш язык — без брокера сообщений, без отдельной базы под очередь.

# Install NATS server
# macOS
brew install nats-server

# Linux
curl -L https://github.com/nats-io/nats-server/releases/download/v2.10.0/nats-server-v2.10.0-linux-amd64.tar.gz | tar xz

# Start
nats-server

# Python client
pip install nats-py

# Node.js client
npm install nats

Публикация задач CAPTCHA в очередь

Издатель (обычно это сам скрапер) делает три вещи за один проход:

  1. Собирает пачку задач с sitekey и pageurl для каждой страницы.
  2. Публикует каждую задачу в subject captcha.tasks отдельным сообщением.
  3. Сбрасывает буфер (flush) и закрывает соединение — публикация в NATS не ждёт подтверждения от воркеров.

Python

import asyncio
import json
import nats


async def publish_captcha_tasks():
    nc = await nats.connect("nats://localhost:4222")

    tasks = [
        {
            "task_id": f"task_{i}",
            "method": "userrecaptcha",
            "sitekey": "6Le-wvkSAAAAAPBMRTvw0Q4Muexq9bi0DJwx_mJ-",
            "pageurl": f"https://example.com/page/{i}"
        }
        for i in range(100)
    ]

    for task in tasks:
        await nc.publish("captcha.tasks", json.dumps(task).encode())
        print(f"Published: {task['task_id']}")

    await nc.flush()
    await nc.close()


asyncio.run(publish_captcha_tasks())

JavaScript

const { connect, StringCodec } = require("nats");

const sc = StringCodec();

async function publishCaptchaTasks() {
  const nc = await connect({ servers: "nats://localhost:4222" });

  for (let i = 0; i < 100; i++) {
    const task = {
      task_id: `task_${i}`,
      method: "userrecaptcha",
      sitekey: "6Le-wvkSAAAAAPBMRTvw0Q4Muexq9bi0DJwx_mJ-",
      pageurl: `https://example.com/page/${i}`,
    };

    nc.publish("captcha.tasks", sc.encode(JSON.stringify(task)));
    console.log(`Published: ${task.task_id}`);
  }

  await nc.flush();
  await nc.close();
}

publishCaptchaTasks();

Воркер: как задачи CAPTCHA попадают в CaptchaAI

Внутри queue group каждое сообщение получает ровно один воркер, даже если параллельно работает несколько экземпляров подписчика на один и тот же subject. На каждой итерации воркер:

  • забирает следующее сообщение из captcha.tasks в составе queue group captcha-workers;
  • отправляет задачу в CaptchaAI (in.php) и опрашивает res.php до готового решения или таймаута;
  • публикует итог в captcha.results, независимо от того, решена задача или завершилась ошибкой.

Python

import asyncio
import json
import os
import nats
import aiohttp

API_KEY = os.environ["CAPTCHAAI_API_KEY"]


async def solve_captcha(session, task):
    """Submit to CaptchaAI and poll for result."""
    # Submit
    async with session.post("https://ocr.captchaai.com/in.php", data={
        "key": API_KEY,
        "method": task["method"],
        "googlekey": task["sitekey"],
        "pageurl": task["pageurl"],
        "json": 1
    }) as resp:
        data = await resp.json(content_type=None)

    if data.get("status") != 1:
        return {"task_id": task["task_id"], "error": data.get("request")}

    captcha_id = data["request"]

    # Poll for result
    for _ in range(60):
        await asyncio.sleep(5)
        async with session.get("https://ocr.captchaai.com/res.php", params={
            "key": API_KEY, "action": "get", "id": captcha_id, "json": 1
        }) as resp:
            result = await resp.json(content_type=None)

        if result.get("status") == 1:
            return {"task_id": task["task_id"], "solution": result["request"]}
        if result.get("request") != "CAPCHA_NOT_READY":
            return {"task_id": task["task_id"], "error": result.get("request")}

    return {"task_id": task["task_id"], "error": "TIMEOUT"}


async def worker(worker_id):
    nc = await nats.connect("nats://localhost:4222")

    # Subscribe with queue group — each message goes to one worker only
    sub = await nc.subscribe("captcha.tasks", queue="captcha-workers")

    print(f"Worker {worker_id} listening...")

    async with aiohttp.ClientSession() as session:
        async for msg in sub.messages:
            task = json.loads(msg.data.decode())
            print(f"Worker {worker_id} processing {task['task_id']}")

            result = await solve_captcha(session, task)

            # Publish result
            await nc.publish(
                "captcha.results",
                json.dumps(result).encode()
            )

            status = "solved" if "solution" in result else result.get("error")
            print(f"  → {task['task_id']}: {status}")


asyncio.run(worker(1))

JavaScript

const { connect, StringCodec } = require("nats");
const axios = require("axios");

const sc = StringCodec();
const API_KEY = process.env.CAPTCHAAI_API_KEY;

function sleep(ms) {
  return new Promise((r) => setTimeout(r, ms));
}

async function solveCaptcha(task) {
  const submitResp = await axios.post(
    "https://ocr.captchaai.com/in.php",
    null,
    {
      params: {
        key: API_KEY,
        method: task.method,
        googlekey: task.sitekey,
        pageurl: task.pageurl,
        json: 1,
      },
    }
  );

  if (submitResp.data.status !== 1) {
    return { task_id: task.task_id, error: submitResp.data.request };
  }

  const captchaId = submitResp.data.request;

  for (let i = 0; i < 60; i++) {
    await sleep(5000);
    const result = await axios.get("https://ocr.captchaai.com/res.php", {
      params: { key: API_KEY, action: "get", id: captchaId, json: 1 },
    });

    if (result.data.status === 1) {
      return { task_id: task.task_id, solution: result.data.request };
    }
    if (result.data.request !== "CAPCHA_NOT_READY") {
      return { task_id: task.task_id, error: result.data.request };
    }
  }

  return { task_id: task.task_id, error: "TIMEOUT" };
}

async function worker(workerId) {
  const nc = await connect({ servers: "nats://localhost:4222" });

  // Queue group subscription — load-balanced across workers
  const sub = nc.subscribe("captcha.tasks", { queue: "captcha-workers" });

  console.log(`Worker ${workerId} listening...`);

  for await (const msg of sub) {
    const task = JSON.parse(sc.decode(msg.data));
    console.log(`Worker ${workerId} processing ${task.task_id}`);

    const result = await solveCaptcha(task);
    nc.publish("captcha.results", sc.encode(JSON.stringify(result)));

    const status = result.solution ? "solved" : result.error;
    console.log(`  → ${task.task_id}: ${status}`);
  }
}

worker(1);

Сбор результатов и статистика решений

async def collect_results():
    nc = await nats.connect("nats://localhost:4222")
    sub = await nc.subscribe("captcha.results")

    solved = 0
    failed = 0

    async for msg in sub.messages:
        result = json.loads(msg.data.decode())

        if "solution" in result:
            solved += 1
            print(f"[SOLVED] {result['task_id']} — {result['solution'][:30]}...")
        else:
            failed += 1
            print(f"[FAILED] {result['task_id']} — {result['error']}")

        print(f"  Stats: {solved} solved, {failed} failed")

asyncio.run(collect_results())

Типичные проблемы и их решения

Большинство проблем с очередью CAPTCHA-задач на NATS сводится к четырём сценариям — все они диагностируются логами воркера и HTTP-монитором сервера, без сторонних инструментов.

Проблема Причина Решение
Сообщения теряются Ядро NATS (pub/sub) не буферизует данные для медленных подписчиков Включите JetStream для сохранения сообщений или увеличьте пропускную способность воркеров
Воркер не получает сообщения Неверное имя queue group или subject Проверьте, что subject и queue group совпадают у издателя и воркера
Сброс соединения Перезапуск NATS-сервера Включите автоматическое переподключение в настройках клиента
Неравномерное распределение Один воркер обрабатывает задачи быстрее других Это нормально — NATS отдаёт сообщения свободным воркерам, более быстрые получают больше

JetStream: когда сообщения нельзя терять

Для задач, потеря которых недопустима — например, если каждая CAPTCHA привязана к платному запросу клиента — включите персистентность JetStream:

async def durable_publisher():
    nc = await nats.connect("nats://localhost:4222")
    js = nc.jetstream()

    # Create stream (one-time setup)
    await js.add_stream(name="CAPTCHA", subjects=["captcha.>"])

    # Publish with acknowledgment
    ack = await js.publish("captcha.tasks", json.dumps(task).encode())
    print(f"Published to stream, seq={ack.seq}")

JetStream добавляет хранение сообщений на диске, возможность повторного воспроизведения потока и доставку без дублей — по сути, надёжность Kafka поверх той же простоты NATS.

Горизонтальное масштабирование воркеров

# Run multiple workers — NATS distributes automatically via queue groups
python worker.py --id=1 &
python worker.py --id=2 &
python worker.py --id=3 &

# Each task goes to exactly one worker
# Add more workers to increase throughput

Добавляя воркеров, вы линейно наращиваете пропускную способность очереди — упирается это не в NATS, а в число потоков вашего плана CaptchaAI и в скорость решения конкретного типа CAPTCHA.

Число воркеров и число занятых потоков CaptchaAI — не одно и то же: воркер может держать несколько задач параллельно (через asyncio/пул соединений), а тариф ограничивает именно количество одновременных потоков, а не число процессов-воркеров.

Частые вопросы

Сколько воркеров нужно для 10 000 задач CAPTCHA в час?

Это ≈2,8 задачи в секунду — сама NATS обрабатывает миллионы сообщений в секунду и узким местом не станет. Реальный лимит — число одновременных потоков вашего плана CaptchaAI и время решения конкретного типа CAPTCHA: считайте от них, а не от пропускной способности очереди.

Чем queue group отличается от обычной подписки NATS?

Обычная подписка на subject доставляет копию каждого сообщения всем подписчикам — так подписан сборщик результатов на captcha.results в примере выше: каждый результат ему нужен целиком. Queue group делает противоположное: все воркеры в одной группе делят между собой один поток сообщений, и каждое сообщение получает только один из них — это и есть балансировка нагрузки из коробки, без внешнего балансировщика или отдельного механизма распределения работы.

Нужен ли JetStream, если воркер иногда падает?

Если задача, потерянная при падении воркера, — это просто повторная отправка без последствий, JetStream не обязателен: core NATS проще и достаточно. Если каждая задача привязана к платному запросу или к строгому SLA, включайте JetStream — он подтверждает доставку и переигрывает недоставленные сообщения.

Как мониторить очередь NATS с задачами CAPTCHA в продакшене?

Считайте разницу между числом опубликованных сообщений в captcha.tasks и результатами в captcha.results — растущий разрыв означает, что воркеры не успевают, и пора добавить ещё один экземпляр worker.py. Встроенный HTTP-монитор NATS-сервера (порт 8222 по умолчанию) отдаёт метрики по подписчикам, очередям и медленным потребителям, которые удобно снимать Prometheus-экспортером и выводить на общий дашборд рядом с метриками самого CaptchaAI (доля успешных решений, время ответа res.php).

Следующие шаги

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