DevOps & Scaling

Messagerie NATS + CaptchaAI : répartition légère des tâches CAPTCHA

Un serveur NATS tient dans un binaire d'une vingtaine de mégaoctets et achemine une tâche CAPTCHA vers un worker disponible en moins d'une milliseconde, sans cluster ni volume à provisionner. Quand vos tâches sont jetables — une résolution perdue se resoumet, elle ne se rejoue pas six mois plus tard — c'est le broker le moins coûteux à exploiter, et l'API CaptchaAI se branche derrière en une trentaine de lignes de Python.

Au programme : publication des tâches, groupe de file d'attente, worker branché sur in.php et res.php, JetStream et dimensionnement des threads.


L'architecture en un coup d'œil

[Scrapers] → Publish → [NATS: captcha.tasks]
                              ↓
                    Queue Group: captcha-workers
                    ├── Worker 1 (solve via CaptchaAI)
                    ├── Worker 2
                    └── Worker 3
                              ↓
                    Publish → [NATS: captcha.results]
                              ↓
                    [Result Subscribers]

Tout tient dans le groupe de file d'attente (« queue group ») : plusieurs workers s'abonnent au sujet captcha.tasks sous le même nom de groupe, et NATS n'envoie chaque message qu'à un seul. Aucun verrouillage à écrire.


NATS, Kafka ou RabbitMQ : lequel pour des tâches CAPTCHA ?

Latences issues des documentations publiées et de mesures observées ; elles varient selon l'environnement et le volume.

Caractéristique NATS Kafka RabbitMQ
Latence < 1 ms 5–10 ms 1–5 ms
Complexité de configuration Binaire unique Cluster + ZooKeeper Modérée
Empreinte mémoire ~20 Mo ~1 Go+ ~200 Mo
Persistance Facultative (JetStream) Intégrée Intégrée
Terrain de prédilection Tâches éphémères, faible latence Diffusion durable Routage complexe

Le critère n'est pas le débit, que les trois encaissent, mais le coût d'exploitation : une tâche CAPTCHA vit quelques dizaines de secondes et n'a aucune valeur historique. Kafka si le flux alimente de l'analytique, RabbitMQ pour du routage fin, NATS sinon.


Installer le serveur et les clients

# Install NATS server
# macOS
brew install nats-server

# Linux
curl -L https://github.com/nats-io/nats-server/releases/download/v2.10.0/nats-server-v2.10.0-linux-amd64.tar.gz | tar xz

# Start
nats-server

# Python client
pip install nats-py

# Node.js client
npm install nats

En production, lancez nats-server en service systemd sur une petite instance, dans la même région que vos workers (OVHcloud à Gravelines, AWS eu-west-3).


Publier les tâches depuis vos scrapers

Le publisher sérialise un dictionnaire en JSON et le pousse sur captcha.tasks. Gardez-y un task_id stable : c'est lui qui recolle le résultat à la page d'origine.

Python

import asyncio
import json
import nats


async def publish_captcha_tasks():
    nc = await nats.connect("nats://localhost:4222")

    tasks = [
        {
            "task_id": f"task_{i}",
            "method": "userrecaptcha",
            "sitekey": "6Le-wvkSAAAAAPBMRTvw0Q4Muexq9bi0DJwx_mJ-",
            "pageurl": f"https://example.com/page/{i}"
        }
        for i in range(100)
    ]

    for task in tasks:
        await nc.publish("captcha.tasks", json.dumps(task).encode())
        print(f"Published: {task['task_id']}")

    await nc.flush()
    await nc.close()


asyncio.run(publish_captcha_tasks())

JavaScript

const { connect, StringCodec } = require("nats");

const sc = StringCodec();

async function publishCaptchaTasks() {
  const nc = await connect({ servers: "nats://localhost:4222" });

  for (let i = 0; i < 100; i++) {
    const task = {
      task_id: `task_${i}`,
      method: "userrecaptcha",
      sitekey: "6Le-wvkSAAAAAPBMRTvw0Q4Muexq9bi0DJwx_mJ-",
      pageurl: `https://example.com/page/${i}`,
    };

    nc.publish("captcha.tasks", sc.encode(JSON.stringify(task)));
    console.log(`Published: ${task.task_id}`);
  }

  await nc.flush();
  await nc.close();
}

publishCaptchaTasks();

Le worker qui résout via CaptchaAI

Chaque worker s'abonne avec l'option queue : ajoutez-en trois, dix ou cinquante, NATS répartit sans configuration. La logique interne est celle de l'API CaptchaAI : un POST sur in.php renvoie un identifiant, puis vous interrogez res.php toutes les 5 s tant que la réponse vaut CAPCHA_NOT_READY.

Deux réglages comptent : l'intervalle de polling — interroger chaque seconde ne raccourcit rien — et le plafond de tentatives, reCAPTCHA v2 étant annoncé sous les 60 s et Turnstile sous les 10 s.

Python

import asyncio
import json
import os
import nats
import aiohttp

API_KEY = os.environ["CAPTCHAAI_API_KEY"]


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

    if data.get("status") != 1:
        return {"task_id": task["task_id"], "error": data.get("request")}

    captcha_id = data["request"]

    # Poll for result
    for _ in range(60):
        await asyncio.sleep(5)
        async with session.get("https://ocr.captchaai.com/res.php", params={
            "key": API_KEY, "action": "get", "id": captcha_id, "json": 1
        }) as resp:
            result = await resp.json(content_type=None)

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

    return {"task_id": task["task_id"], "error": "TIMEOUT"}


async def worker(worker_id):
    nc = await nats.connect("nats://localhost:4222")

    # Subscribe with queue group — each message goes to one worker only
    sub = await nc.subscribe("captcha.tasks", queue="captcha-workers")

    print(f"Worker {worker_id} listening...")

    async with aiohttp.ClientSession() as session:
        async for msg in sub.messages:
            task = json.loads(msg.data.decode())
            print(f"Worker {worker_id} processing {task['task_id']}")

            result = await solve_captcha(session, task)

            # Publish result
            await nc.publish(
                "captcha.results",
                json.dumps(result).encode()
            )

            status = "solved" if "solution" in result else result.get("error")
            print(f"  → {task['task_id']}: {status}")


asyncio.run(worker(1))

JavaScript

const { connect, StringCodec } = require("nats");
const axios = require("axios");

const sc = StringCodec();
const API_KEY = process.env.CAPTCHAAI_API_KEY;

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

  return { task_id: task.task_id, error: "TIMEOUT" };
}

async function worker(workerId) {
  const nc = await connect({ servers: "nats://localhost:4222" });

  // Queue group subscription — load-balanced across workers
  const sub = nc.subscribe("captcha.tasks", { queue: "captcha-workers" });

  console.log(`Worker ${workerId} listening...`);

  for await (const msg of sub) {
    const task = JSON.parse(sc.decode(msg.data));
    console.log(`Worker ${workerId} processing ${task.task_id}`);

    const result = await solveCaptcha(task);
    nc.publish("captcha.results", sc.encode(JSON.stringify(result)));

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

worker(1);

Collecter les résultats

async def collect_results():
    nc = await nats.connect("nats://localhost:4222")
    sub = await nc.subscribe("captcha.results")

    solved = 0
    failed = 0

    async for msg in sub.messages:
        result = json.loads(msg.data.decode())

        if "solution" in result:
            solved += 1
            print(f"[SOLVED] {result['task_id']} — {result['solution'][:30]}...")
        else:
            failed += 1
            print(f"[FAILED] {result['task_id']} — {result['error']}")

        print(f"  Stats: {solved} solved, {failed} failed")

asyncio.run(collect_results())

Ce collecteur compte, mais c'est surtout là que vous brancherez vos métriques : taux de réussite, temps de résolution médian, codes d'erreur. Un ERROR_ZERO_BALANCE et un pic de TIMEOUT n'appellent pas la même réponse.


Ajouter JetStream quand une tâche perdue coûte cher

Le cœur de NATS ne conserve rien : si aucun worker n'écoute au moment de la publication, le message est perdu. Acceptable pour un scraping relançable, pas pour un formulaire à soumettre une seule fois — activez alors JetStream.

async def durable_publisher():
    nc = await nats.connect("nats://localhost:4222")
    js = nc.jetstream()

    # Create stream (one-time setup)
    await js.add_stream(name="CAPTCHA", subjects=["captcha.>"])

    # Publish with acknowledgment
    ack = await js.publish("captcha.tasks", json.dumps(task).encode())
    print(f"Published to stream, seq={ack.seq}")

JetStream ajoute la persistance sur disque, la relecture et l'accusé de réception côté publisher : une fiabilité de livraison proche de Kafka avec un seul binaire à exploiter. En contrepartie, surveillez le stockage et fixez une rétention sur le stream CAPTCHA.


Combien de workers lancer ? Alignez-les sur vos threads

# Run multiple workers — NATS distributes automatically via queue groups
python worker.py --id=1 &
python worker.py --id=2 &
python worker.py --id=3 &

# Each task goes to exactly one worker
# Add more workers to increase throughput

Multiplier les workers ne sert à rien si le plan CaptchaAI ne suit pas. La facturation repose sur les threads — une résolution en vol, résolutions illimitées par thread : BASIC ($15/mois, 5 threads), STANDARD ($30/mois, 15 threads), ADVANCE ($90/mois, 50 threads), PREMIUM ($170/mois, 100 threads), jusqu'à VIP-3 ($7,500/mois, 5 000 threads).

La règle : les résolutions simultanées de tous vos workers ne doivent pas dépasser les threads du plan, soit 15 processus mono-tâche au plus sur STANDARD. Un thread enchaînant bien plus de Turnstile (moins de 10 s) que de reCAPTCHA v2 (moins de 60 s), dimensionnez d'après votre mix de types.


Exemple : veille tarifaire pour un distributeur belge

Une équipe data surveille chaque nuit les prix de quinze concurrents : 4 000 pages parcourues, dont 300 déclenchent un défi CAPTCHA. Trois briques suffisent — un nats-server sur une VM à 2 vCPU, huit workers Python dans l'image Docker du crawler, un collecteur qui écrit dans PostgreSQL. Le plan STANDARD ($30/mois, 15 threads) couvre la fenêtre nocturne et le pic se lisse tout seul.

À noter : vos messages circulent en clair sur le réseau interne. Pas d'identifiants clients ni de données personnelles dans le payload — un task_id opaque et une URL suffisent, et vos obligations RGPD de minimisation s'en trouvent simplifiées.


Dépannage

Problème Cause probable Correctif
Des tâches disparaissent sous forte charge Le cœur pub/sub ne tamponne pas les consommateurs lents (« slow consumer ») Basculer captcha.tasks sur JetStream ou ajouter des workers
Un worker ne reçoit rien Sujet ou nom de groupe différent de celui du publisher Aligner captcha.tasks et queue="captcha-workers" des deux côtés
Coupures de connexion répétées Redémarrage du serveur ou réseau instable Garder la reconnexion automatique active et le traitement idempotent via le task_id
Beaucoup de TIMEOUT après des CAPCHA_NOT_READY Plus de résolutions en vol que de threads sur le plan Réduire le nombre de workers ou monter de palier
Répartition inégale entre workers Certains traitent des types plus rapides Normal : NATS sert le worker disponible ; comparez les débits par type

FAQ

Combien de workers NATS puis-je faire tourner avec mon plan CaptchaAI ?

Autant que de threads inclus dans le plan si chaque worker traite une tâche à la fois : 5 sur BASIC ($15/mois, 5 threads), 15 sur STANDARD, 50 sur ADVANCE. Au-delà, les tâches patientent.

Que se passe-t-il si un worker s'arrête pendant l'interrogation du résultat ?

Le message a déjà été consommé et le cœur de NATS ne le redistribue pas. Avec JetStream et un accusé de réception envoyé après publication du résultat, la tâche repart vers un autre worker.

Faut-il un sujet NATS distinct par type de CAPTCHA ?

Souvent oui : captcha.tasks.recaptcha et captcha.tasks.turnstile autorisent des groupes de workers séparés, dimensionnés par type. Le wildcard captcha.> couvre la supervision.

Puis-je router des tâches hCaptcha ou FunCaptcha vers ces workers ?

Non — hCaptcha et FunCaptcha (Arkose Labs) ne sont pas pris en charge. Les types couverts sont reCAPTCHA v2 et v3, Cloudflare Turnstile, Cloudflare Challenge, GeeTest v3 et les CAPTCHA image ou grille ; GeeTest v4 est annoncé comme à venir, et CaptchaFox, Friendly Captcha et Lemin restent en bêta.


Prochaines étapes

Montez le serveur, lancez trois workers, mesurez le débit — créez votre clé API CaptchaAI pour commencer.

Guides associés :

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