DevOps y Escalado

RabbitMQ + CaptchaAI: integración de cola de mensajes

Un worker se cae en mitad de una resolución. ¿Pierdes la tarea? Con RabbitMQ no: el mensaje sigue en la cola hasta que alguien lo confirma, y otro worker lo retoma en segundos. Esa es la razón práctica por la que un pipeline de CAPTCHA en producción acaba montado sobre un broker y no sobre una lista en memoria del propio proceso.

Aquí montas esa integración con Python y pika: colas duraderas, confirmación manual, dead letter exchange y enrutamiento por tipo de CAPTCHA, contra la API de CaptchaAI (in.php para enviar, res.php para consultar el resultado).


Qué te aporta RabbitMQ en un pipeline de CAPTCHA

Capacidad Qué te resuelve en producción
Colas duraderas Las tareas sobreviven a un reinicio del broker
Confirmación manual (ack) Si el worker muere, la tarea vuelve a la cola
Dead letter exchange Lo que falla de forma repetida se aparta para revisarlo
Colas con prioridad Un CAPTCHA de un flujo interactivo pasa antes que uno de lote
Routing keys Cada tipo de CAPTCHA llega a workers especializados

La pieza que más se nota es la confirmación manual: mientras el worker no llama a basic_ack, el broker considera la tarea en vuelo y la redistribuye si el proceso muere. Esa red de seguridad evita huecos silenciosos en tus datos.


RabbitMQ o Redis: cómo elegir sin complicarte

Si tu caso es "una cola, unos cuantos workers y baja latencia", una lista de Redis te sirve. RabbitMQ compensa cuando necesitas repartir por tipo, aislar los fallos en una cola aparte, usar prioridades o auditar desde el panel del broker por qué una tarea concreta nunca devolvió token.


Preparar el broker y el cliente de Python

Levanta RabbitMQ con la imagen de gestión (panel en el puerto 15672) e instala las dependencias:

# Docker
docker run -d --hostname rabbitmq \
  -p 5672:5672 -p 15672:15672 \
  rabbitmq:3-management

# Python client
pip install pika requests

Paso 1 — Productor: publicar tareas en la cola

El productor declara toda la topología antes de publicar nada: el dead letter exchange, la cola de fallidos, la cola principal con TTL y la cola de resultados. Declarar es idempotente, así que puedes arrancar varios productores sin miedo. Fíjate en delivery_mode=2: sin ese marcado de persistencia, una cola duradera sigue perdiendo su contenido al reiniciar el broker.

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()

Paso 2 — Consumidor: el worker que resuelve

El consumidor toma una tarea, la envía a CaptchaAI, sondea el resultado y publica el token en la cola de resultados. La clave está en prefetch_count=1: cada worker retiene una sola tarea, así que una resolución lenta no bloquea al resto del grupo. Cuando algo falla, responde con basic_nack y requeue=False, y el mensaje sale por el dead letter exchange hacia captcha.failed en lugar de repetirse en bucle.

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()

Paso 3 — Recolector de resultados

El recolector lee la cola de resultados y devuelve los tokens agrupados por identificador de tarea. Es el punto natural para enganchar tu lógica: guardar el token, reenviarlo al formulario o registrar el tiempo de resolución.

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

Enrutamiento por tipo de CAPTCHA

Cuando mezclas tipos, conviene que cada uno tenga su propia cola. Así puedes dar más workers a reCAPTCHA v2, mantener Cloudflare Turnstile en una ruta rápida y aislar el OCR de imágenes, que suele tener un perfil de tiempos distinto:

# 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),
    )

También facilita las pruebas: puedes parar los workers de un tipo sin frenar los demás.


Cuántos workers ejecutar (y cómo se relaciona con tu plan)

El límite real no es RabbitMQ, sino cuántas resoluciones tienes en vuelo a la vez. CaptchaAI factura por thread concurrente, con resoluciones ilimitadas dentro de cada thread: BASIC ($15/mes) incluye 5 threads, STANDARD ($30/mes) 15, ADVANCE ($90/mes) 50 y PREMIUM ($170/mes) 100. Los precios se mantienen en USD.

De ahí sale la regla: el número de consumidores con prefetch_count=1 no debería superar los threads de tu plan. Con 40 workers y 15 threads, las peticiones sobrantes esperan turno y tu latencia media empeora. Mide la longitud de captcha.tasks en hora punta y ajusta primero el plan, después los workers.


Escenario: monitorización de portales públicos y marketplaces

En España y Latinoamérica es habitual el equipo que vigila portales protegidos con CAPTCHA: páginas de cita previa, trámites administrativos o fichas de producto en marketplaces regionales. El trabajo llega a ráfagas —cientos de comprobaciones por la mañana y flujo plano el resto del día—, que es justo donde una cola en medio se nota.

El productor publica y sigue, sin saber nada de tu capacidad de resolución. Los workers absorben la ráfaga al ritmo que permitan tus threads y lo que falla queda en captcha.failed con su mensaje original, para ver si el problema fue el sitekey, la URL o un tiempo de espera. Antes de automatizar cualquier portal, revisa sus términos de servicio y la normativa de protección de datos aplicable.


Errores frecuentes y cómo corregirlos

Síntoma Causa probable Corrección
Se pierden mensajes al reiniciar Cola no duradera o mensajes sin persistir Declara la cola con durable=True y publica con delivery_mode=2
Un worker se queda "atascado" Resolución larga con demasiadas tareas reservadas Mantén prefetch_count=1 en cada consumidor
La cola de fallidos crece sin parar Fallos persistentes en los parámetros de la tarea Revisa los mensajes de captcha.failed y corrige sitekey o pageurl
Se cae la conexión con el broker Tiempo de espera del heartbeat Define el intervalo de heartbeat y añade lógica de reconexión

Preguntas frecuentes

¿RabbitMQ añade latencia frente a llamar a la API directamente?

Muy poca: el salto por el broker se mide en milisegundos y la resolución en segundos. Lo que ganas a cambio —reintentos, aislamiento de fallos y desacople— compensa de sobra en cualquier flujo asíncrono.

¿Qué pasa con las tareas en vuelo si el broker se reinicia?

Los mensajes persistidos en colas duraderas se recuperan al arrancar y los que un worker tenía sin confirmar se reparten de nuevo. El peor caso es una resolución repetida, no una tarea perdida.

¿Necesito colas con prioridad?

Solo si conviven flujos interactivos y de lote. Si hay un usuario esperando ante un formulario, publica esa tarea con priority alta y deja el trabajo en segundo plano en el nivel por defecto; con una sola clase de carga, la prioridad solo añade complejidad.

¿Cómo evito que la cola de mensajes fallidos se descontrole?

Ponle una alerta por longitud y revísala a diario. El patrón suele repetirse: un sitekey caducado, una URL mal formada o un tipo sin configurar. Corrige el origen y republica los mensajes en captcha.tasks.

¿Puedo usar la misma cola para reCAPTCHA v2 y Cloudflare Turnstile?

Sí, siempre que cada mensaje lleve su method y sus parámetros. Aun así, separarlos con routing keys te da métricas por tipo y te permite dimensionar los workers según el tiempo de resolución de cada uno.


Guías relacionadas


Colas fiables de principio a fin: empieza con CaptchaAI y RabbitMQ.

Los comentarios están deshabilitados para este artículo.