DevOps y Escalado

Apache Kafka + CaptchaAI: transmisión de procesamiento de tareas CAPTCHA

Si tu pipeline de scraping ya genera más CAPTCHA de los que una cola sencilla puede absorber, Apache Kafka es la pieza que te falta. Kafka desacopla el envío de tareas de su resolución con un flujo de mensajes duradero y ordenado, y tus scrapers no se bloquean esperando un token. En esta guía conectas Kafka con la API de CaptchaAI para resolver miles de CAPTCHA por hora con workers escalables.

Piensa en una agencia de datos en Buenos Aires que monitorea precios en marketplaces regionales: en un pico de lanzamientos genera decenas de miles de páginas con reCAPTCHA en pocas horas. Donde Redis o RabbitMQ se saturan, Kafka acumula las tareas en un topic duradero y los workers las consumen a su ritmo, sin perder ninguna.

Kafka encaja cuando necesitas:

  • durabilidad y reproducción de mensajes,
  • alto rendimiento, del orden de 100.000 mensajes/seg,
  • escalado horizontal mediante grupos de consumidores.

Para volúmenes por debajo de 1.000 tareas por hora, Redis o RabbitMQ bastan; a partir de ahí, Kafka rinde mejor.

Arquitectura del pipeline: dos topics de Kafka

[Scrapers] → Produce → [Kafka: captcha-tasks topic]
                              ↓
                    [CAPTCHA Worker Group]
                    (consume tasks, solve via CaptchaAI)
                              ↓
                    Produce → [Kafka: captcha-results topic]
                              ↓
                    [Result Consumer Group]
                    (process solutions, update database)

Usa dos topics de Kafka para separar responsabilidades:

  • captcha-tasks: parámetros de CAPTCHA a la espera de resolución.
  • captcha-results: tokens ya resueltos, listos para su uso posterior.

Requisitos previos

Antes de empezar necesitas:

  • un broker de Kafka accesible en localhost:9092 (o la dirección de tu clúster),
  • las librerías cliente de Python o Node.js,
  • tu CAPTCHAAI_API_KEY exportada en el entorno.

Instala las dependencias del cliente:

# Python
pip install kafka-python requests

# Node.js
npm install kafkajs axios

Paso 1: crear los topics de Kafka

kafka-topics.sh --create --topic captcha-tasks \
  --partitions 6 --replication-factor 1 \
  --bootstrap-server localhost:9092

kafka-topics.sh --create --topic captcha-results \
  --partitions 6 --replication-factor 1 \
  --bootstrap-server localhost:9092

Seis particiones permiten hasta seis consumidores en paralelo por grupo y fijan tu techo de concurrencia.

Paso 2: productor de tareas (lado del scraper)

Cada mensaje lleva los campos que el worker necesita:

  • task_id: identificador único de la tarea,
  • method: tipo de CAPTCHA a resolver,
  • sitekey y pageurl: parámetros del desafío.

Usar el task_id como clave (key) mantiene cada tarea en la misma partición y preserva su orden.

Python

import json
from kafka import KafkaProducer

producer = KafkaProducer(
    bootstrap_servers=["localhost:9092"],
    value_serializer=lambda v: json.dumps(v).encode("utf-8"),
    key_serializer=lambda k: k.encode("utf-8") if k else None,
    acks="all",  # Wait for all replicas to confirm
    retries=3
)


def enqueue_captcha(task_id, sitekey, pageurl, captcha_type="userrecaptcha"):
    """Send a CAPTCHA task to Kafka."""
    task = {
        "task_id": task_id,
        "method": captcha_type,
        "sitekey": sitekey,
        "pageurl": pageurl,
        "submitted_at": __import__("time").time()
    }

    future = producer.send(
        "captcha-tasks",
        key=task_id,  # Key ensures same task goes to same partition
        value=task
    )
    future.get(timeout=10)  # Block until confirmed
    return task_id


# Submit tasks
enqueue_captcha("task_001", "6Le-wvkSAAAAAPBMRTvw0Q4Muexq9bi0DJwx_mJ-", "https://example.com")
enqueue_captcha("task_002", "6Le-wvkSAAAAAPBMRTvw0Q4Muexq9bi0DJwx_mJ-", "https://example.com")
producer.flush()

JavaScript

const { Kafka } = require("kafkajs");

const kafka = new Kafka({
  clientId: "captcha-producer",
  brokers: ["localhost:9092"],
});

const producer = kafka.producer();

async function enqueueCaptcha(taskId, sitekey, pageurl) {
  await producer.connect();

  const task = {
    task_id: taskId,
    method: "userrecaptcha",
    sitekey: sitekey,
    pageurl: pageurl,
    submitted_at: Date.now(),
  };

  await producer.send({
    topic: "captcha-tasks",
    messages: [{ key: taskId, value: JSON.stringify(task) }],
  });
}

(async () => {
  await enqueueCaptcha(
    "task_001",
    "6Le-wvkSAAAAAPBMRTvw0Q4Muexq9bi0DJwx_mJ-",
    "https://example.com"
  );
  await producer.disconnect();
})();

Validación mínima antes de publicar en Kafka

REQUIRED_TASK_FIELDS = {"task_id", "method", "sitekey", "pageurl"}


def validate_task(task: dict) -> dict:
    missing = REQUIRED_TASK_FIELDS - task.keys()
    if missing:
        raise ValueError(f"Missing task fields: {sorted(missing)}")
    return task

Validar los campos obligatorios antes de producir evita que un mensaje mal formado desperdicie un thread de resolución en el worker.

Paso 3: worker de CAPTCHA (consumidor y solver)

En cada iteración, el worker:

  • consume una tarea de captcha-tasks,
  • la envía a CaptchaAI vía in.php y sondea res.php hasta obtener el token,
  • publica el resultado en captcha-results y confirma el offset.

El commit manual del offset (después de procesar, no antes) evita perder tareas si un worker se cae. El campo method decide qué se resuelve: userrecaptcha (reCAPTCHA v2 y v3), turnstile (Cloudflare Turnstile) o geetest (GeeTest v3).

Python

import json
import os
import time
import requests
from kafka import KafkaConsumer, KafkaProducer

API_KEY = os.environ["CAPTCHAAI_API_KEY"]

consumer = KafkaConsumer(
    "captcha-tasks",
    bootstrap_servers=["localhost:9092"],
    group_id="captcha-workers",
    value_deserializer=lambda m: json.loads(m.decode("utf-8")),
    auto_offset_reset="earliest",
    enable_auto_commit=False,  # Manual commit after processing
    max_poll_records=10
)

result_producer = KafkaProducer(
    bootstrap_servers=["localhost:9092"],
    value_serializer=lambda v: json.dumps(v).encode("utf-8")
)


def solve_captcha(task):
    """Submit to CaptchaAI and poll for result."""
    # Submit
    resp = requests.post("https://ocr.captchaai.com/in.php", data={
        "key": API_KEY,
        "method": task["method"],
        "googlekey": task["sitekey"],
        "pageurl": task["pageurl"],
        "json": 1
    })
    data = resp.json()

    if data.get("status") != 1:
        return {"error": data.get("request")}

    captcha_id = data["request"]

    # Poll for result
    for _ in range(60):
        time.sleep(5)
        result = requests.get("https://ocr.captchaai.com/res.php", params={
            "key": API_KEY,
            "action": "get",
            "id": captcha_id,
            "json": 1
        }).json()

        if result.get("status") == 1:
            return {"solution": result["request"]}
        if result.get("request") != "CAPCHA_NOT_READY":
            return {"error": result.get("request")}

    return {"error": "TIMEOUT"}


# Main consumer loop
print("CAPTCHA worker started. Waiting for tasks...")
for message in consumer:
    task = message.value
    print(f"Processing {task['task_id']}...")

    result = solve_captcha(task)
    result["task_id"] = task["task_id"]
    result["solved_at"] = time.time()

    # Publish result
    result_producer.send("captcha-results", value=result)
    result_producer.flush()

    # Commit offset after successful processing
    consumer.commit()
    print(f"  → {task['task_id']}: {'solved' if 'solution' in result else result.get('error')}")

JavaScript

const { Kafka } = require("kafkajs");
const axios = require("axios");

const API_KEY = process.env.CAPTCHAAI_API_KEY;

const kafka = new Kafka({
  clientId: "captcha-worker",
  brokers: ["localhost:9092"],
});

const consumer = kafka.consumer({ groupId: "captcha-workers" });
const producer = kafka.producer();

function sleep(ms) {
  return new Promise((resolve) => setTimeout(resolve, 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 { 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 { solution: result.data.request };
    if (result.data.request !== "CAPCHA_NOT_READY")
      return { error: result.data.request };
  }

  return { error: "TIMEOUT" };
}

async function run() {
  await consumer.connect();
  await producer.connect();
  await consumer.subscribe({ topic: "captcha-tasks", fromBeginning: false });

  await consumer.run({
    eachMessage: async ({ message }) => {
      const task = JSON.parse(message.value.toString());
      console.log(`Processing ${task.task_id}...`);

      const result = await solveCaptcha(task);
      result.task_id = task.task_id;
      result.solved_at = Date.now();

      await producer.send({
        topic: "captcha-results",
        messages: [{ value: JSON.stringify(result) }],
      });

      console.log(
        `  → ${task.task_id}: ${result.solution ? "solved" : result.error}`
      );
    },
  });
}

run();

Escalar los workers de resolución con grupos de consumidores

Los grupos de consumidores de Kafka reparten automáticamente las particiones entre los workers:

# 6 partitions, 3 workers → each worker gets 2 partitions
Worker-1: partitions 0, 1
Worker-2: partitions 2, 3
Worker-3: partitions 4, 5

# Add Worker-4 → rebalance
Worker-1: partitions 0, 1
Worker-2: partitions 2
Worker-3: partitions 3, 4
Worker-4: partition 5

Escala hasta el número de particiones; más allá de ese punto, añade más particiones para sumar workers. Al reajustar el grupo, ten en cuenta que:

  • un rebalanceo pausa brevemente el consumo en todos los workers,
  • tener más particiones que workers deja particiones ociosas,
  • todos los workers deben compartir el mismo group_id.

Solución de problemas frecuentes

Antes de mirar métricas, descarta las causas más comunes cuando el pipeline se atasca:

Problema Causa Solución
El consumer lag no para de crecer Los workers no dan abasto con la tasa de tareas Añade más instancias de workers (hasta el número de particiones)
Resultados duplicados El worker se cae antes de confirmar el offset Añade una verificación de idempotencia sobre task_id en el consumidor de resultados
Rebalanceos demasiado frecuentes Workers que se caen o se reinician Aumenta session.timeout.ms; revisa si hay OOM
Tareas mal repartidas Distribución de claves deficiente Usa claves aleatorias o más particiones

Nota: al extraer datos de forma automatizada, respeta los términos de servicio del sitio y la normativa de protección de datos aplicable.

Monitoreo del consumer lag en Kafka

La señal clave de si tu pipeline va al día es el consumer lag. Consúltalo así:

kafka-consumer-groups.sh --describe --group captcha-workers \
  --bootstrap-server localhost:9092
Métrica Saludable Advertencia
Consumer lag < 100 > 1000 (añade workers)
Mensajes/seg entrantes Coincide con la tasa del scraper Los picos indican ráfagas
Mensajes/seg salientes Coincide con la tasa de entrada Quedarse atrás = cuello de botella

Preguntas frecuentes

¿Cuántas particiones necesito para mi volumen?

Como cada partición admite un worker activo por grupo, las particiones marcan el techo de tu paralelismo. Como guía práctica:

  • 6–12 particiones cubren varios miles de tareas por hora,
  • súbelas antes de saturarte, porque añadir es sencillo y reducir no,
  • nunca pongas más workers que particiones en un mismo grupo.

¿Qué plan de CaptchaAI necesito para alimentar varios workers en paralelo?

CaptchaAI factura por thread concurrente (no por resolución), con resoluciones ilimitadas por thread. Alinea los threads con tu concurrencia:

  • BASIC ($15/mes, 5 threads): pruebas y volúmenes bajos.
  • ADVANCE ($90/mes, 50 threads): producción media.
  • PREMIUM ($170/mes, 100 threads): producción sostenida con varios workers.

¿Cómo evito resolver el mismo CAPTCHA dos veces?

Apóyate en dos mecanismos complementarios:

  • Commit manual del offset: el worker confirma solo después de publicar el resultado.
  • Clave idempotente: una verificación sobre task_id en el consumidor de resultados descarta cualquier duplicado.

Para un CAPTCHA sin solución, usa un límite de reintentos y deriva el mensaje a un topic captcha-dead-letter en lugar de bloquear la partición.

¿Puedo mezclar distintos tipos de CAPTCHA en el mismo topic?

Sí. El campo method le indica al worker qué resolver, así que reCAPTCHA v2/v3, Cloudflare Turnstile y GeeTest v3 conviven en captcha-tasks sin topics separados por tipo.

Artículos relacionados

  • Procesamiento por lotes de resultados CAPTCHA en streaming

Próximos pasos

Monta tu pipeline de CAPTCHA en streaming: obtén tu API key de CaptchaAI y conecta Kafka para alto rendimiento.

Guías relacionadas:

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