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_KEYexportada 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,sitekeyypageurl: 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.phpy sondeares.phphasta obtener el token, - publica el resultado en
captcha-resultsy 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_iden 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:
- Cola de Redis para procesamiento distribuido
- Integración con la cola de mensajes RabbitMQ
- Cómo resolver 10.000 tareas por hora