Kafka sert à découpler la soumission des tâches CAPTCHA de leur résolution : vos producteurs empilent les tâches au rythme du scraping, un groupe de workers les consomme au sien, et aucune rafale de trafic ne fait tomber le pipeline. C'est l'architecture à sortir dès que vous dépassez quelques milliers de résolutions par heure, que plusieurs services lisent les résultats, ou que vous devez rejouer un flux après incident.
Architecture du pipeline
Le principe tient en deux topics et deux groupes de consommateurs. Les scrapers publient les tâches, des workers les résolvent via CaptchaAI, puis republient les tokens dans un second topic que d'autres services consomment à leur rythme.
[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)
Deux topics séparent les responsabilités ; chaque côté évolue sans casser l'autre :
captcha-tasks– les paramètres CAPTCHA en attente de résolution ;captcha-results– les tokens résolus, prêts à être utilisés en aval.
Faut-il vraiment Kafka pour vos CAPTCHA ?
Kafka n'est pas gratuit à exploiter : cluster, supervision, équipe qui sait le faire tourner. Adoptez-le si l'un de ces points s'applique :
- plusieurs producteurs et un débit continu élevé ;
- un besoin de rejouer ou historiser les flux de tâches ;
- plusieurs services qui consomment les résultats en aval ;
- une équipe déjà à l'aise avec l'exploitation d'un cluster Kafka.
Si la colonne de droite vous décrit mieux, restez sur une file plus simple.
| Signal | Kafka a du sens | Une file plus simple suffit |
|---|---|---|
| Volume de tâches | Plusieurs producteurs, gros débit continu | Quelques workers, trafic modéré |
| Besoin de relecture | Il faut rejouer et historiser les flux | Les tâches sont éphémères et peu critiques |
| Consommateurs en aval | Plusieurs services lisent les résultats | Un seul pipeline applicatif traite tout |
| Coût d'exploitation | L'équipe sait déjà exploiter Kafka | Le poids opérationnel dépasserait le gain |
Prérequis
# Python
pip install kafka-python requests
# Node.js
npm install kafkajs axios
Il vous faut par ailleurs :
- un broker Kafka accessible sur
localhost:9092(ou l'adresse de votre cluster) ; - une clé API CaptchaAI exposée dans la variable d'environnement
CAPTCHAAI_API_KEY.
Étape 1 : créer les topics
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
Six partitions autorisent jusqu'à six consommateurs en parallèle par groupe. Le nombre de partitions plafonne le parallélisme : prévoyez-en large dès le départ.
Étape 2 : le producteur de tâches (côté scraper)
Le producteur envoie une tâche par CAPTCHA rencontré. La clé de message (task_id) fixe la partition cible ; avec acks="all", la publication n'est confirmée qu'une fois répliquée, donc aucune tâche ne se perd si un broker tombe.
| Champ | Rôle |
|---|---|
task_id |
Identifiant unique, qui sert aussi de clé de partition |
method |
Type de CAPTCHA (userrecaptcha, turnstile, geetest…) |
sitekey |
Clé publique du site cible |
pageurl |
URL de la page qui affiche le défi |
submitted_at |
Horodatage de mise en file |
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();
})();
Étape 3 : le worker CAPTCHA (consommateur + solveur)
Le worker consomme une tâche, la résout via CaptchaAI, puis republie le token dans captcha-results.
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();
Trois détails d'exploitation :
- l'offset n'est validé qu'après la publication du résultat (
enable_auto_commit=False), donc une tâche interrompue est rejouée plutôt que perdue ; - la résolution interroge le résultat toutes les 5 secondes, jusqu'à 60 fois ;
- un échec renvoie une erreur au lieu de bloquer la partition.
Mettre les workers à l'échelle
Les groupes de consommateurs Kafka répartissent automatiquement les partitions entre les workers d'un même groupe :
# 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
Augmentez le nombre de workers jusqu'au nombre de partitions ; au-delà, ils restent inactifs. Pour aller plus loin, ajoutez d'abord des partitions, puis des workers.
Dimensionner les threads CaptchaAI
Chaque worker occupe un thread CaptchaAI pendant qu'il interroge un résultat, puis le libère. Dimensionnez votre plan en trois temps :
- comptez vos workers actifs en pointe : c'est votre besoin en threads simultanés ;
- réservez autant de threads CaptchaAI — 6 workers demandent au moins 6 threads ;
- choisissez le plan qui couvre ce chiffre : STANDARD ($30/mois, 15 threads) ici, ADVANCE ($90/mois, 50 threads) pour un topic bien plus large, BASIC ($15/mois, 5 threads) pour un prototype.
Hébergé en Europe, un cluster Kafka managé chez OVHcloud ou Scaleway — ou des workers en région eu-west-3 (Paris) — garde la latence basse. Côté conformité, journalisez le task_id et le statut, jamais les données personnelles récupérées : c'est la base pour rester aligné avec le RGPD.
Superviser le pipeline
Le principal indicateur de santé est le lag des consommateurs — l'écart entre les tâches publiées et les tâches traitées :
kafka-consumer-groups.sh --describe --group captcha-workers \
--bootstrap-server localhost:9092
| Métrique | Sain | Alerte |
|---|---|---|
| Lag du consommateur | < 100 | > 1000 (ajoutez des workers) |
| Messages/s entrants | Suit le rythme du scraping | Les pics signalent une rafale |
| Messages/s sortants | Suit le débit entrant | Décrochage = goulet d'étranglement |
Un lag qui grimpe signifie que la résolution ne suit pas la production : ajoutez des workers, puis des partitions au-delà.
Dépannage
| Problème | Cause | Correctif |
|---|---|---|
| Le lag du consommateur s'accentue | Les workers ne suivent pas le rythme des tâches | Ajoutez des workers (jusqu'au nombre de partitions) |
| Résultats en double | Le worker plante avant de valider l'offset | Ajoutez un contrôle d'idempotence sur task_id côté consommateur de résultats |
| Rééquilibrage trop fréquent | Workers qui plantent ou redémarrent | Augmentez session.timeout.ms et vérifiez les OOM (out of memory) |
| Tâches mal réparties | Mauvaise distribution des clés | Utilisez des clés variées ou davantage de partitions |
FAQ
Combien de threads CaptchaAI faut-il pour mon nombre de workers ?
Autant que de workers actifs en pointe : chaque worker occupe un thread pendant qu'il interroge un résultat. Six workers demandent au moins six threads, couverts par le plan STANDARD ($30/mois, 15 threads).
Comment éviter qu'une tâche soit traitée deux fois ?
Ne validez l'offset qu'après avoir publié le résultat, et ajoutez un contrôle d'idempotence sur task_id côté consommateur. Si un worker plante avant le commit, la tâche est rejouée ; le contrôle écarte alors le doublon.
Pourquoi Kafka plutôt que Redis ou RabbitMQ ?
Pour la durabilité et la relecture : Kafka conserve les messages, laisse rejouer un flux entier, encaisse un gros débit et répartit la charge via les groupes de consommateurs. En dessous d'environ 1 000 tâches/heure et sans besoin de replay, Redis ou RabbitMQ restent plus simples.
Que faire d'un CAPTCHA impossible à résoudre ?
Fixez une limite de tentatives dans le worker, puis publiez la tâche sur un topic captcha-dead-letter pour inspection manuelle. Ne bloquez jamais une partition avec des tentatives infinies.
Articles connexes
- Traiter les résultats CAPTCHA par lots en streaming
Prochaines étapes
Lancez vos pipelines CAPTCHA en streaming : récupérez votre clé API CaptchaAI et branchez Kafka.
- Une file Redis pour le traitement distribué
- Intégrer une file de messages RabbitMQ
- Résoudre 10 000 tâches par heure