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
- Cola de tareas con Redis para procesamiento distribuido
- Resolver CAPTCHA por lotes con varias tareas en paralelo
Colas fiables de principio a fin: empieza con CaptchaAI y RabbitMQ.