DevOps & Scaling

RabbitMQ + CaptchaAI: интеграция очереди сообщений

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.

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