Tutoriales

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

Un lote de 500 tareas CAPTCHA se completa de manera desigual: algunas se resuelven en 8 segundos, otras tardan 45. Esperar a que finalice cada tarea antes de procesar los resultados hace perder el tiempo entre la primera y la última solución. La transmisión permite que su canal descendente consuma cada resultado en el momento en que llega.

Streaming versus proceso por lotes y luego

Enfoque Tiempo hasta el primer resultado Memoria Latencia de canalización
espera por todos Después de la tarea más lenta Todos los resultados en la memoria. Alto
Transmitir como resuelto Después de la tarea más rápida Un resultado a la vez Bajo
Microlote (trozos de 10) Después del primer trozo 10 resultados a la vez Medio

Python: generador asíncrono para resultados de transmisión

Usando asyncio y aiohttp, cada solución produce inmediatamente a través de un generador asíncrono:

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

Instalar dependencias:

pip install aiohttp

JavaScript: patrón de transmisión EventEmitter

Node.js utiliza un enfoque basado en eventos: emite cada resultado a medida que se resuelve:

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

Cuándo utilizar Streaming frente a Collect-All

Escenario Enfoque
Envíos de formularios utilizando tokens Stream: envíe cada formulario tan pronto como llegue el token
Exportación CSV de todos los resultados Recopilar todo: escribir una vez cuando se complete el lote
Panel con progreso en vivo Stream: actualiza la interfaz de usuario en cada evento de resultado
Lote con dependencias entre tareas Recopilar todo: procesar en orden una vez finalizado
Lotes grandes (más de 1000) Stream: reduce el uso máximo de memoria

Solución de problemas

Problema causa Solución
Los resultados llegan en orden aleatorio. Normal: el streaming produce el rendimiento más rápido primero Utilice result.index para volver a asignar la tarea original
La memoria sigue creciendo durante la transmisión Almacenar todos los resultados en una matriz Procesar y descartar resultados en el controlador.
El primer resultado tarda demasiado Todas las tareas enviadas simultáneamente Envíos escalonados con semáforo o límite de concurrencia
Advertencia de EventEmitter: MaxListenersExceeded Hay demasiados oyentes en la transmisión Utilice setMaxListeners() o asegúrese de que haya un oyente por tipo de evento
El generador asíncrono se bloquea Tarea no resuelta en conjunto pendiente Agregue tiempo de espera a poll_task; asegurar que todos los futuros estén completos o con errores

Preguntas frecuentes

¿La transmisión aumenta las llamadas API en comparación con el lote?

No: se produce la misma cantidad de llamadas de envío y consulta de cualquier manera. La transmisión solo cambia cuando su aplicación procesa cada resultado, no cuántas llamadas API se realizan.

¿Cómo mantengo el orden de las tareas durante la transmisión?

Cada resultado lleva su index original. Si el orden es importante para el procesamiento posterior, el búfer da como resultado una estructura ordenada y vacía ejecuciones contiguas (como el reensamblaje de paquetes TCP).

¿Puedo combinar streaming con checkpointing?

Sí. Agregue cada resultado a un archivo de punto de control a medida que llegue. Al reanudar, cargue el punto de control, filtre los índices completados y vuelva a procesar solo las tareas restantes.

Artículos relacionados

Próximos pasos

Procese las soluciones CAPTCHA en el momento en que llegan:obtenga su clave API CaptchaAIy construir canales de streaming.

Guías relacionadas:

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