DevOps y Escalado

Mensajería NATS + CaptchaAI: distribución ligera de tareas CAPTCHA

¿Necesitas repartir miles de resoluciones de CAPTCHA entre varios workers sin montar un clúster de Kafka? NATS está pensado para eso: un broker de un solo binario, sin JVM y sin persistencia en disco por defecto, con latencia inferior a un milisegundo. Reparte las tareas entre tus workers de CaptchaAI con una huella mínima.

Cuándo NATS es la elección correcta

Una tarea CAPTCHA es efímera: si se pierde una, la reenvías y listo, sin llevar un historial de cada mensaje. Eso cambia qué broker conviene. NATS encaja cuando:

  • El trabajo es de corta vida y reenviar una tarea perdida no cuesta nada.
  • Quieres latencia por debajo del milisegundo entre publicar la tarea y que un worker la reciba.
  • Prefieres un solo binario que arranca en segundos frente a un clúster con dependencias que operar.
  • El volumen sube y baja cada hora, y quieres sumar o quitar workers sin reconfigurar nada.

Si en cambio necesitas un registro duradero de cada evento para auditoría o reprocesamiento masivo, ahí Kafka sigue siendo la herramienta indicada.

Cómo fluyen las tareas de un extremo a otro

[Scrapers] → Publish → [NATS: captcha.tasks]
                              ↓
                    Queue Group: captcha-workers
                    ├── Worker 1 (solve via CaptchaAI)
                    ├── Worker 2
                    └── Worker 3
                              ↓
                    Publish → [NATS: captcha.results]
                              ↓
                    [Result Subscribers]

El recorrido tiene cuatro etapas claras:

  • Los scrapers publican cada tarea en el subject captcha.tasks.
  • Un grupo de cola llamado captcha-workers reparte esas tareas entre los workers suscritos.
  • Cada worker resuelve su CAPTCHA con CaptchaAI y publica el resultado en captcha.results.
  • Un suscriptor aparte recoge los resultados y lleva la cuenta.

La pieza clave son los grupos de cola: NATS entrega cada mensaje a un único worker del grupo, de modo que no resuelves —ni pagas— el mismo CAPTCHA por duplicado.

Instalación y requisitos previos

Instala el servidor NATS, arráncalo y añade el cliente de tu lenguaje.

# 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

Publicar tareas desde el scraper

Cada scraper serializa la tarea en JSON y la publica en el subject captcha.tasks.

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

El worker que resuelve el CAPTCHA

Al compartir el mismo grupo de cola, NATS trata a los workers como un pool y entrega cada tarea a uno solo. Dentro del worker, el flujo con CaptchaAI sigue tres pasos:

  • Envías la tarea a in.php y recibes un ID de captcha.
  • Sondeas res.php hasta que la resolución esté lista.
  • Publicas la solución (o el error) en 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);

Recolectar los resultados

Los workers publican cada resolución en captcha.results, donde un suscriptor aparte lleva la cuenta de éxitos y fallos.

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

JetStream cuando las tareas no pueden perderse

El pub/sub básico es de tipo "dispara y olvida". Si necesitas que ninguna tarea se pierda, activa la persistencia de 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 añade persistencia en disco, reproducción y entrega exactamente una vez: como Kafka, pero conservando la simplicidad de NATS.

Escalar los workers horizontalmente

# 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

Arrancas más procesos con el mismo grupo de cola y NATS les reparte tareas al instante, sin tocar nada más. El techo no lo pone NATS —que mueve millones de mensajes por segundo—, sino cuántas resoluciones sostienes en paralelo con CaptchaAI.

Cuántos threads de CaptchaAI dimensionar

Aquí está el matiz que confunde a muchos equipos: sumar workers de NATS no acelera nada por sí solo. Lo que marca el ritmo es tu número de threads en CaptchaAI, porque cada thread es un CAPTCHA en vuelo.

  • El límite real es el paralelismo de resolución, no cuántos procesos worker levantes.
  • El plan BASIC ($15/mes, 5 threads) da un arranque holgado para tres o cuatro workers.
  • STANDARD ($30/mes, 15 threads) sube el paralelismo cuando el volumen crece.
  • Ese costo mensual predecible en USD es cómodo para equipos que facturan en monedas volátiles.

Piensa en un ejemplo cercano: una agencia que monitorea precios en marketplaces regionales —un MercadoLibre o un Amazon.es— puede arrancar con un plan modesto y crecer por threads a medida que suma catálogos, sin rehacer la arquitectura de mensajería.

NATS frente a Kafka y RabbitMQ

Si aún dudas entre brokers, esta tabla de referencia resume por qué NATS gana en pipelines de CAPTCHA de baja latencia:

Característica NATS Kafka RabbitMQ
Latencia < 1 ms 5-10 ms 1-5 ms
Complejidad de configuración Un solo binario Clúster + ZooKeeper Moderada
Consumo de memoria ~20 MB ~1 GB+ ~200 MB
Persistencia Opcional (JetStream) Integrada Integrada
Ideal para Tareas efímeras, baja latencia Streaming duradero Enrutamiento complejo

Resolución de problemas

Problema Causa Solución
Se pierden mensajes El pub/sub básico de NATS no almacena en búfer para consumidores lentos Usa JetStream para persistencia o aumenta la capacidad de consumo
El worker no recibe mensajes Nombre de grupo de cola o subject incorrecto Verifica que el subject y el grupo de cola coincidan con el publicador
Se reinicia la conexión El servidor NATS se reinició Activa la reconexión automática en las opciones del cliente
Distribución desigual Un worker procesa más rápido que otros Normal: NATS reparte entre los workers disponibles; los más rápidos reciben más

Preguntas frecuentes

¿Los grupos de cola de NATS evitan que dos workers resuelvan el mismo CAPTCHA?

Sí. Mientras todos compartan el mismo grupo de cola sobre el mismo subject, NATS entrega cada mensaje a uno solo. No hay riesgo de resolver un CAPTCHA por duplicado ni de gastar dos veces el mismo thread.

¿Cuántos threads de CaptchaAI necesito para mis workers NATS?

Tantos como resoluciones simultáneas quieras sostener. Cada thread es un CAPTCHA en vuelo: da igual cuántos workers levantes, si solo tienes 5 threads resolverás 5 a la vez. El plan BASIC ($15/mes) incluye 5 threads y STANDARD ($30/mes) sube a 15; dimensiona según tu paralelismo real, no según el número de workers.

¿Puedo combinar workers de Python y JavaScript en el mismo grupo de cola?

Sí. NATS no distingue el lenguaje del suscriptor: mientras compartan el subject captcha.tasks, el grupo de cola captcha-workers y el mismo formato JSON en el mensaje, puedes mezclar workers de Python y JavaScript en el mismo pool sin problema.

¿Qué ocurre si el servidor NATS se reinicia a mitad de un lote?

Con el pub/sub básico, los mensajes en vuelo se pierden y basta con reenviarlos. Si eso no es aceptable, activa JetStream: persiste las tareas en disco y las reproduce tras el reinicio. En ambos casos, habilita la reconexión automática del cliente.

Empieza a distribuir tus tareas CAPTCHA

Levanta un par de workers y deja que los grupos de cola hagan el reparto. Obtén tu API key de CaptchaAI y resuelve tu primer lote en minutos.

Guías relacionadas:

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