DevOps & Scaling

Apache Kafka + CaptchaAI : traitement des tâches CAPTCHA en streaming

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 :

  1. comptez vos workers actifs en pointe : c'est votre besoin en threads simultanés ;
  2. réservez autant de threads CaptchaAI — 6 workers demandent au moins 6 threads ;
  3. 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.

Les commentaires sont désactivés pour cet article.