Если нужно раскидать тысячи задач 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 в очередь
Издатель (обычно это сам скрапер) делает три вещи за один проход:
- Собирает пачку задач с
sitekeyиpageurlдля каждой страницы. - Публикует каждую задачу в subject
captcha.tasksотдельным сообщением. - Сбрасывает буфер (
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 groupcaptcha-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).