DevOps & Mise à l'Échelle

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

NATS est un système de messagerie léger et hautes performances : pas de JVM, pas de persistance de disque par défaut, latence inférieure à la milliseconde. Pour la distribution de tâches CAPTCHA où vous avez besoin de rapidité et de simplicité par rapport à la durabilité de Kafka, NATS est la solution idéale.

Pourquoi NATS pour les tâches CAPTCHA

Caractéristique NATS Kafka LapinMQ
Latence < 1 ms 5-10 ms 1-5 ms
Complexité de configuration Binaire simple Cluster + ZooKeeper Modéré
Empreinte mémoire ~20 Mo ~1 Go+ ~200 Mo
Persistance Facultatif (JetStream) Intégré Intégré
Idéal pour Tâches éphémères, faible latence Diffusion durable Routage complexe

Les tâches CAPTCHA sont éphémères : si une tâche est perdue, vous la soumettez à nouveau. La simplicité et la rapidité de NATS en font un choix naturel.

Architecture

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

Les groupes de files d'attente NATS distribuent automatiquement les messages entre les travailleurs : chaque tâche est confiée à exactement un travailleur.

Conditions préalables

# 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

Éditeur de tâches (Scraper)

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();

Travailleur CAPTCHA (abonné au groupe de file d'attente)

Les groupes de files d'attente garantissent que chaque message est envoyé à exactement un travailleur, même si plusieurs travailleurs sont en cours d'exécution.

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);

Collecteur de 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())

NATS JetStream pour la durabilité

Pour les tâches qui ne doivent pas être perdues, activez la persistance 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 du disque, la capacité de relecture et la livraison exactement une fois – similaire à Kafka mais avec la simplicité de NATS.

Travailleurs à l'échelle

# 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

Dépannage

Problème Parce que Corriger
Messages abandonnés Le noyau pub/sub de NATS ne met pas en mémoire tampon les consommateurs lents Utilisez JetStream pour la persistance ou augmentez la capacité des consommateurs
Le travailleur ne reçoit pas de messages Nom ou sujet du groupe de file d'attente incorrect Vérifier que le sujet et le groupe de file d'attente correspondent à l'éditeur
Réinitialisation de la connexion Redémarrage du serveur NATS Activer la reconnexion automatique dans les options client
Répartition inégale Un travailleur traite plus rapidement que les autres Normal : NATS distribue aux travailleurs disponibles ; les travailleurs les plus rapides obtiennent plus

FAQ

Quand dois-je utiliser NATS au lieu de Redis ou Kafka ?

Utilisez NATS lorsque vous souhaitez une infrastructure minimale (un seul binaire, aucune dépendance), une latence inférieure à la milliseconde et que vous n'avez pas besoin d'un stockage persistant des messages. Utilisez Kafka pour un streaming durable, Redis pour la mise en cache + le combo pub/sub.

NATS peut-il gérer plus de 10 000 tâches CAPTCHA par heure ?

Facilement. NATS gère des millions de messages par seconde. Le goulot d'étranglement sera le temps de résolution de CaptchaAI, et non le débit NATS.

Ai-je besoin de JetStream ?

Seulement si vous avez besoin de persévérance. Pour la plupart des flux de travail CAPTCHA, le NATS de base est suffisant : si une tâche est perdue, vous pouvez la soumettre à nouveau. Activez JetStream pour les pistes d'audit ou les exigences de traitement unique.

Prochaines étapes

Distribuez des tâches CAPTCHA avec NATS -récupérez votre clé API CaptchaAIet faire tourner les travailleurs légers.

Guides associés :

  • Intégration du streaming Kafka
  • File d'attente Redis pour le traitement distribué
  • Intégration de la file d'attente de messages RabbitMQ
Les commentaires sont désactivés pour cet article.