Tutoriales

Transmisión de resultados por lotes: procesamiento de soluciones CAPTCHA a medida que llegan

Respuesta corta: no esperes al lote completo. Consume cada token al resolverse y tu pipeline arranca con la primera respuesta, no con la última. Las tareas nunca terminan a la vez —unas caen en 8 segundos y otras tardan 45—, así que bloquearte hasta la más lenta regala minutos.

  • El patrón en Python con asyncio y generadores asíncronos.
  • El equivalente en Node.js con EventEmitter.
  • Cuándo transmitir y cuándo recopilar todo.

Todo con la API de CaptchaAI: envío a in.php, sondeo en res.php.

Por qué transmitir gana tiempo frente a esperar el lote

La diferencia no está en la API, sino en cuándo tu código lee el resultado:

Enfoque Primer resultado Memoria Latencia
Esperar a todas las tareas Tras la más lenta Todo en memoria Alta
Transmitir al resolverse Tras la más rápida Un resultado Baja
Microlotes (bloques de 10) Tras el primer bloque 10 resultados Media

El streaming no acelera la resolución, pero recupera el hueco entre la primera solución y la última.

Ejemplo: una suite de QA nocturna que no quiere esperar

Un estudio en Bogotá corre cada noche pruebas de humo sobre los formularios de alta de sus clientes: unos 300 casos detrás de un CAPTCHA. Recopilando todo, si el despliegue rompió el login esa señal llega media hora tarde.

Transmitiendo, cada caso se marca al llegar su token y el equipo aborta tras tres fallos seguidos. Cortar antes no ahorra dinero —se factura por thread, no por resolución—, pero adelanta el diagnóstico:

  • BASIC ($15/mes, 5 threads) cubre una suite pequeña diaria.
  • ADVANCE ($90/mes, 50 threads) entra cuando la ventana nocturna se acorta.
  • Respeta los términos de servicio y la normativa de protección de datos vigente.

Python: generador asíncrono que emite cada resultado

La clave es asyncio.wait(..., return_when=FIRST_COMPLETED): devuelve el control al terminar la primera tarea, no todas. El semáforo respeta tus threads.

import asyncio
import aiohttp
import time

API_KEY = "YOUR_API_KEY"
SUBMIT_URL = "https://ocr.captchaai.com/in.php"
RESULT_URL = "https://ocr.captchaai.com/res.php"


async def submit_task(session, task_data):
    """Submit a single CAPTCHA task."""
    params = {
        "key": API_KEY,
        "method": task_data.get("method", "userrecaptcha"),
        "json": 1,
    }
    if params["method"] == "userrecaptcha":
        params["googlekey"] = task_data["sitekey"]
        params["pageurl"] = task_data["pageurl"]
    elif params["method"] == "turnstile":
        params["sitekey"] = task_data["sitekey"]
        params["pageurl"] = task_data["pageurl"]

    async with session.post(SUBMIT_URL, data=params) as resp:
        result = await resp.json(content_type=None)
        if result.get("status") != 1:
            return None, result.get("request", "unknown")
        return result["request"], None


async def poll_task(session, task_id, timeout=300):
    """Poll until solved or timeout."""
    start = time.monotonic()
    while time.monotonic() - start < timeout:
        await asyncio.sleep(5)
        params = {"key": API_KEY, "action": "get", "id": task_id, "json": 1}
        async with session.get(RESULT_URL, params=params) as resp:
            result = await resp.json(content_type=None)

        if result.get("request") == "CAPCHA_NOT_READY":
            continue
        if result.get("status") == 1:
            return result["request"], None
        return None, result.get("request", "unknown")

    return None, "TIMEOUT"


async def solve_one(session, index, task_data, semaphore):
    """Solve a single task within concurrency limits."""
    async with semaphore:
        start = time.monotonic()
        task_id, error = await submit_task(session, task_data)
        if error:
            return {"index": index, "status": "failed", "error": error, "time": 0}

        token, error = await poll_task(session, task_id)
        elapsed = time.monotonic() - start

        if token:
            return {"index": index, "status": "solved", "token": token, "time": round(elapsed, 1)}
        return {"index": index, "status": "failed", "error": error, "time": round(elapsed, 1)}


async def stream_results(tasks, max_concurrent=20):
    """
    Async generator that yields each result as it completes.
    Results arrive in completion order, not submission order.
    """
    semaphore = asyncio.Semaphore(max_concurrent)

    async with aiohttp.ClientSession() as session:
        pending = set()
        for i, task in enumerate(tasks):
            coro = solve_one(session, i, task, semaphore)
            pending.add(asyncio.ensure_future(coro))

        while pending:
            done, pending = await asyncio.wait(pending, return_when=asyncio.FIRST_COMPLETED)
            for future in done:
                yield future.result()


async def main():
    tasks = [
        {"sitekey": "SITE_KEY", "pageurl": f"https://example.com/page{i}"}
        for i in range(50)
    ]

    solved = 0
    failed = 0

    async for result in stream_results(tasks, max_concurrent=15):
        # Process each result immediately
        if result["status"] == "solved":
            solved += 1
            print(f"  [{solved + failed}/{len(tasks)}] Task {result['index']} SOLVED in {result['time']}s")

            # Use token immediately — don't wait for batch
            # await submit_form(result["token"])
            # await save_to_database(result)
        else:
            failed += 1
            print(f"  [{solved + failed}/{len(tasks)}] Task {result['index']} FAILED: {result['error']}")

    print(f"\nDone: {solved} solved, {failed} failed")


asyncio.run(main())

Instala la dependencia:

pip install aiohttp

Los resultados llegan en orden de finalización, no de envío, y cada uno arrastra su index hasta la tarea de origen.

Node.js: patrón EventEmitter para transmitir soluciones

El equivalente en Node.js es un emisor de eventos: publica un result por tarea cerrada y un done final. El consumidor solo escucha y actúa.

const { EventEmitter } = require("events");

const API_KEY = "YOUR_API_KEY";
const SUBMIT_URL = "https://ocr.captchaai.com/in.php";
const RESULT_URL = "https://ocr.captchaai.com/res.php";

class CaptchaStream extends EventEmitter {
  constructor(maxConcurrent = 15) {
    super();
    this.maxConcurrent = maxConcurrent;
    this.active = 0;
    this.queue = [];
    this.total = 0;
    this.completed = 0;
  }

  async submitAndPoll(index, taskData) {
    const params = new URLSearchParams({
      key: API_KEY,
      method: taskData.method || "userrecaptcha",
      googlekey: taskData.sitekey,
      pageurl: taskData.pageurl,
      json: "1",
    });

    const start = Date.now();
    const submitResp = await (await fetch(SUBMIT_URL, { method: "POST", body: params })).json();

    if (submitResp.status !== 1) {
      return { index, status: "failed", error: submitResp.request, time: 0 };
    }

    const taskId = submitResp.request;
    for (let i = 0; i < 60; i++) {
      await new Promise((r) => setTimeout(r, 5000));
      const url = `${RESULT_URL}?key=${API_KEY}&action=get&id=${taskId}&json=1`;
      const poll = await (await fetch(url)).json();

      if (poll.request === "CAPCHA_NOT_READY") continue;
      const elapsed = ((Date.now() - start) / 1000).toFixed(1);
      if (poll.status === 1) return { index, status: "solved", token: poll.request, time: elapsed };
      return { index, status: "failed", error: poll.request, time: elapsed };
    }
    return { index, status: "failed", error: "TIMEOUT", time: ((Date.now() - start) / 1000).toFixed(1) };
  }

  async processNext() {
    if (this.queue.length === 0 || this.active >= this.maxConcurrent) return;

    const { index, taskData } = this.queue.shift();
    this.active++;

    try {
      const result = await this.submitAndPoll(index, taskData);
      this.emit("result", result);
    } catch (err) {
      this.emit("result", { index, status: "failed", error: err.message });
    } finally {
      this.active--;
      this.completed++;

      if (this.completed === this.total) {
        this.emit("done");
      } else {
        this.processNext();
      }
    }
  }

  start(tasks) {
    this.total = tasks.length;
    this.queue = tasks.map((taskData, index) => ({ index, taskData }));

    // Launch initial batch
    const initial = Math.min(this.maxConcurrent, tasks.length);
    for (let i = 0; i < initial; i++) {
      this.processNext();
    }
    return this;
  }
}

// Usage
const tasks = Array.from({ length: 50 }, (_, i) => ({
  sitekey: "SITE_KEY",
  pageurl: `https://example.com/page${i}`,
}));

const stream = new CaptchaStream(15);
let solved = 0, failed = 0;

stream.on("result", (result) => {
  if (result.status === "solved") {
    solved++;
    console.log(`[${solved + failed}/${tasks.length}] Task ${result.index} SOLVED (${result.time}s)`);
    // Use token immediately
    // submitForm(result.token);
  } else {
    failed++;
    console.log(`[${solved + failed}/${tasks.length}] Task ${result.index} FAILED: ${result.error}`);
  }
});

stream.on("done", () => {
  console.log(`\nComplete: ${solved} solved, ${failed} failed`);
});

stream.start(tasks);

El bucle espera 5 segundos entre consultas: sondear más a menudo no acelera nada y gasta llamadas.

Cuándo transmitir y cuándo recopilar todo

  • Formularios que usan el token: transmite y envía cada uno al llegar.
  • Exportación CSV del lote: recopila y escribe una sola vez al final.
  • Panel con progreso en vivo: transmite y actualiza por evento.
  • Dependencias entre tareas: recopila y procesa en orden.
  • Lotes de más de 1.000 tareas: transmite para reducir el pico de memoria.

Regla práctica: si el consumidor trabaja con resultados sueltos, transmite; si necesita el conjunto ordenado, recopila.

Problemas frecuentes y cómo resolverlos

Problema Causa Solución
Resultados desordenados Normal: primero lo más rápido Reasigna con result.index
La memoria crece sin parar Acumulas todo en un array Procesa y descarta en el handler
El primer resultado tarda Enviaste todo a la vez Escalona con el semáforo
Aviso MaxListenersExceeded Demasiados listeners Uno por evento o setMaxListeners()
El generador se cuelga Tarea atascada en pendientes Añade timeout en poll_task

Si los fallos se repiten, revisa primero tu concurrencia.

Preguntas frecuentes

¿El streaming consume más threads de mi plan?

No. Los threads los ocupan las tareas en vuelo, y ese número lo fija max_concurrent.

¿Sirve este patrón para todos los tipos que resuelve CaptchaAI?

Sí, el flujo de envío y sondeo es idéntico; solo cambian los parámetros:

  • reCAPTCHA v2 y v3, Cloudflare Turnstile y Challenge, GeeTest v3, imagen/OCR.
  • CaptchaFox, Friendly Captcha y Lemin, disponibles en beta.

¿Con qué frecuencia debo consultar el resultado?

Cada 5 segundos, como en los ejemplos: bajar el intervalo multiplica las llamadas a res.php sin adelantar nada.

¿Qué pasa si una tarea del lote falla?

Llega igual, con status: "failed" y su código de error. Trátalo en el mismo handler: registra, decide si reintentas y deja correr el lote.

Artículos relacionados

Próximos pasos

Consume cada solución en cuanto llega: obtén tu clave API de CaptchaAI y monta tu pipeline de streaming. Sigue por aquí:

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