Redis решает очередь CAPTCHA-задач быстро, но при падении воркера необработанные сообщения из памяти пропадают без следа. RabbitMQ закрывает именно эту проблему: постоянные (durable) очереди переживают перезапуск брокера, подтверждение сообщений (ack/nack) не даёт задаче исчезнуть при сбое воркера, а dead letter exchange уводит неудачные попытки в отдельную очередь для разбора, а не роняет их молча. Ниже — рабочая интеграция: producer, который ставит CAPTCHA-задачи в очередь, consumer-воркер, который решает их через API CaptchaAI, и маршрутизация по типу CAPTCHA для профильных обработчиков.
Что RabbitMQ даёт очереди CAPTCHA-задач
| Особенность | Что это даёт |
|---|---|
| Durable-очереди | Задачи переживают перезапуск брокера |
| Подтверждение сообщений (ack/nack) | Ни одна задача не теряется при сбое воркера |
| Dead letter exchange | Неудачные задачи уходят на разбор, а не пропадают молча |
| Приоритетные очереди | Срочные CAPTCHA обрабатываются вне очереди |
| Routing key | Задача уходит к воркеру, который умеет решать именно этот тип CAPTCHA |
Подготовка окружения
Поднимите брокер и поставьте клиентскую библиотеку:
# Docker
docker run -d --hostname rabbitmq \
-p 5672:5672 -p 15672:15672 \
rabbitmq:3-management
# Python client
pip install pika requests
Порт 15672 открывает Management UI — по нему удобно смотреть глубину очереди и статистику consumer'ов, не заходя в код.
Producer: постановка задач в очередь
Producer объявляет очереди при первом запуске и публикует задачи с флагом delivery_mode=2 (persistent), чтобы сообщение пережило перезапуск брокера. Основная очередь captcha.tasks настроена с dead-letter-маршрутизацией и TTL в 5 минут — если задача не была обработана за это время, RabbitMQ сам перекладывает её в captcha.failed вместо того, чтобы держать очередь забитой.
import json
import uuid
import pika
class CaptchaProducer:
"""Submit CAPTCHA tasks to RabbitMQ."""
def __init__(self, rabbitmq_url="amqp://guest:guest@localhost:5672/"):
self.connection = pika.BlockingConnection(
pika.URLParameters(rabbitmq_url),
)
self.channel = self.connection.channel()
self._setup_queues()
def _setup_queues(self):
"""Declare durable queues and exchanges."""
# Dead letter exchange for failed tasks
self.channel.exchange_declare(
exchange="captcha.dlx",
exchange_type="direct",
durable=True,
)
self.channel.queue_declare(
queue="captcha.failed",
durable=True,
)
self.channel.queue_bind(
queue="captcha.failed",
exchange="captcha.dlx",
routing_key="failed",
)
# Main task queue with dead letter routing
self.channel.queue_declare(
queue="captcha.tasks",
durable=True,
arguments={
"x-dead-letter-exchange": "captcha.dlx",
"x-dead-letter-routing-key": "failed",
"x-message-ttl": 300000, # 5 min TTL
},
)
# Results queue
self.channel.queue_declare(
queue="captcha.results",
durable=True,
)
def submit(self, method, params, priority=0):
"""Submit a CAPTCHA task."""
task_id = str(uuid.uuid4())[:8]
task = {
"id": task_id,
"method": method,
"params": params,
}
self.channel.basic_publish(
exchange="",
routing_key="captcha.tasks",
body=json.dumps(task),
properties=pika.BasicProperties(
delivery_mode=2, # Persistent
priority=priority,
message_id=task_id,
),
)
return task_id
def close(self):
self.connection.close()
# Usage
producer = CaptchaProducer()
task_id = producer.submit("userrecaptcha", {
"googlekey": "SITE_KEY",
"pageurl": "https://example.com",
}, priority=5)
print(f"Submitted: {task_id}")
producer.close()
Consumer: воркер, который решает CAPTCHA
Consumer забирает задачи по одной (prefetch_count=1) — это ключевая настройка: без неё RabbitMQ может выдать воркеру пачку сообщений сразу, и одна медленная reCAPTCHA v2 заблокирует обработку остальных. Внутри _solve — стандартный цикл CaptchaAI: POST на in.php, затем опрос res.php каждые 5 секунд до готовности токена или таймаута. При успехе воркер публикует результат в captcha.results и подтверждает сообщение (basic_ack); при исключении — отклоняет его (basic_nack, requeue=False), и RabbitMQ сам перекладывает задачу в dead letter exchange.
import json
import os
import time
import pika
import requests
class CaptchaConsumer:
"""RabbitMQ consumer that solves CAPTCHAs."""
def __init__(self, api_key, rabbitmq_url="amqp://guest:guest@localhost:5672/"):
self.api_key = api_key
self.base = "https://ocr.captchaai.com"
self.connection = pika.BlockingConnection(
pika.URLParameters(rabbitmq_url),
)
self.channel = self.connection.channel()
# Process one task at a time
self.channel.basic_qos(prefetch_count=1)
def start(self):
"""Start consuming tasks."""
self.channel.basic_consume(
queue="captcha.tasks",
on_message_callback=self._handle_task,
)
print("Worker started. Waiting for tasks...")
self.channel.start_consuming()
def _handle_task(self, ch, method, properties, body):
"""Process a single CAPTCHA task."""
task = json.loads(body)
task_id = task["id"]
print(f"Processing {task_id}...")
try:
token = self._solve(task["method"], task["params"])
# Publish result
result = {
"task_id": task_id,
"status": "success",
"token": token,
}
ch.basic_publish(
exchange="",
routing_key="captcha.results",
body=json.dumps(result),
properties=pika.BasicProperties(delivery_mode=2),
)
# Acknowledge message (remove from queue)
ch.basic_ack(delivery_tag=method.delivery_tag)
print(f"{task_id} solved successfully")
except Exception as e:
print(f"{task_id} failed: {e}")
# Reject and send to dead letter queue
ch.basic_nack(
delivery_tag=method.delivery_tag,
requeue=False, # Goes to DLX
)
def _solve(self, captcha_method, params, timeout=120):
resp = requests.post(f"{self.base}/in.php", data={
"key": self.api_key,
"method": captcha_method,
"json": 1,
**params,
}, 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")
# Run worker
if __name__ == "__main__":
consumer = CaptchaConsumer(os.environ["CAPTCHAAI_KEY"])
consumer.start()
Если брокер и воркеры разнесены по разным дата-центрам, держите RabbitMQ и обработчики в одном регионе — это стабилизирует heartbeat и снижает число ложных Connection reset, особенно на нестабильных мобильных или спутниковых каналах у удалённых команд.
Сборщик результатов
Отдельный процесс вычитывает captcha.results и собирает готовые токены по task_id, пока не наберётся нужное количество или не истечёт таймаут:
import json
import pika
class ResultCollector:
"""Collect task results from the results queue."""
def __init__(self, rabbitmq_url="amqp://guest:guest@localhost:5672/"):
self.connection = pika.BlockingConnection(
pika.URLParameters(rabbitmq_url),
)
self.channel = self.connection.channel()
self.results = {}
def collect(self, expected_count, timeout=120):
"""Collect a specific number of results."""
deadline = time.time() + timeout
while len(self.results) < expected_count and time.time() < deadline:
method, _, body = self.channel.basic_get(
queue="captcha.results",
auto_ack=True,
)
if body:
result = json.loads(body)
self.results[result["task_id"]] = result
time.sleep(0.5)
return self.results
Такая схема отделяет постановку задач от чтения результатов — producer не блокируется в ожидании ответа, а сборщик может жить в отдельном процессе или даже отдельном сервисе.
Маршрутизация задач по типу CAPTCHA
Если в очереди смешаны reCAPTCHA, Turnstile и обычные image-CAPTCHA, направляйте каждый тип своему воркеру через direct-exchange с routing key — так медленная reCAPTCHA v2 не будет держать в очереди быстрые задачи по Turnstile:
# Setup exchanges and queues
channel.exchange_declare(
exchange="captcha.types",
exchange_type="direct",
durable=True,
)
# Queue per type
for captcha_type in ["recaptcha", "turnstile", "image"]:
channel.queue_declare(queue=f"captcha.{captcha_type}", durable=True)
channel.queue_bind(
queue=f"captcha.{captcha_type}",
exchange="captcha.types",
routing_key=captcha_type,
)
# Submit with routing
def submit_routed(channel, captcha_type, task):
channel.basic_publish(
exchange="captcha.types",
routing_key=captcha_type,
body=json.dumps(task),
properties=pika.BasicProperties(delivery_mode=2),
)
Сколько воркеров и какой план CaptchaAI нужен
Пропускную способность ограничивают не ядра CPU, а SLA-скорость решения на стороне CaptchaAI. reCAPTCHA v2 решается менее чем за 60 секунд, Cloudflare Turnstile — менее чем за 10 секунд. При prefetch_count=1 каждый воркер держит один поток CaptchaAI занятым: четыре воркера на очереди из reCAPTCHA v2 — это до четырёх параллельных решений одновременно, а не количество ядер сервера. Для такой нагрузки достаточно STANDARD ($30/мес, 15 потоков) с запасом, а для смешанной очереди с бóльшим пиком — ADVANCE ($90/мес, 50 потоков). Считать количество воркеров имеет смысл от целевого потока задач в час, а не наоборот.
Типичные проблемы и их решение
| Проблема | Причина | Решение |
|---|---|---|
| Сообщения теряются при падении брокера | Очередь объявлена без durable=True |
Установите durable=True и delivery_mode=2 |
| Воркер завис на одной задаче | Долгое решение CAPTCHA держит весь prefetch | Установите prefetch_count=1 на каждый воркер |
| Dead letter очередь растёт | Задачи стабильно проваливаются | Разберите содержимое captcha.failed и проверьте параметры запроса |
| Соединение обрывается | Не настроен heartbeat | Задайте интервал heartbeat и добавьте логику переподключения |
Часто задаваемые вопросы
Сколько потоков CaptchaAI нужно под очередь RabbitMQ?
Отталкивайтесь от целевого числа задач в час и SLA-скорости типа CAPTCHA, а не от количества воркеров. Для очереди из reCAPTCHA v2 (до 60 секунд на решение) STANDARD (15 потоков) обычно перекрывает 4–8 параллельных воркеров с запасом; для смешанной нагрузки с пиками берите ADVANCE (50 потоков).
Что делать, если RabbitMQ-соединение регулярно рвётся у удалённой команды?
Чаще всего дело в нестабильной сети (мобильный интернет, слабый Wi-Fi) и отсутствии heartbeat. Задайте разумный интервал heartbeat в клиенте, добавьте логику переподключения с backoff и по возможности держите брокер в том же облачном регионе, что и воркеры.
Можно ли автоматически повторять неудачные задачи?
Да. Настройте retry-exchange с задержкой через TTL: сообщения, отклонённые воркером (basic_nack), уходят туда и через заданное время возвращаются в основную очередь без ручного вмешательства.
Как посмотреть глубину очереди и статистику воркеров без кода?
Management UI RabbitMQ на порту 15672 показывает длину каждой очереди, количество consumer'ов и скорость обработки в реальном времени — обычно этого достаточно, чтобы понять, где затор: в producer, в очереди или в воркере.
RabbitMQ или Redis для очереди CAPTCHA-задач?
Берите RabbitMQ, когда важны durable-очереди, dead letter exchange или маршрутизация по типу CAPTCHA. Redis подойдёт для более простой и низколатентной очереди, если потеря отдельной задачи при сбое не критична.
Связанные руководства
Надёжная очередь для CAPTCHA-задач — начните с CaptchaAI и RabbitMQ.