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