DevOps & Scaling

File d'attente Redis + CaptchaAI : traitement CAPTCHA distribué

Un script qui envoie ses CAPTCHA un par un et attend chaque token plafonne à quelques dizaines de résolutions par heure, quel que soit le plan souscrit. Le correctif tient en une ligne d'architecture : placez une file d'attente Redis entre vos producteurs et vos workers, et laissez plusieurs processus la consommer en parallèle. Redis joue trois rôles à la fois : tampon des tâches, registre des résultats et canal de notification. Ce guide monte le pipeline complet en Python — gestionnaire de file, workers, producteur, écoute Pub/Sub, dimensionnement et supervision.


Pourquoi Redis pour répartir vos résolutions CAPTCHA

Résoudre un CAPTCHA est une opération d'attente, pas de calcul. Un token reCAPTCHA v2 revient en moins de 60 secondes, un Cloudflare Turnstile en moins de 10 secondes : pendant tout ce temps, votre processus attend sans consommer de CPU. Paralléliser cette attente est le seul levier de débit réel.

Redis suffit largement : une liste pour la file, un hash pour les résultats, un canal Pub/Sub pour les notifications — sur un composant que vos équipes exploitent souvent déjà. Il découple surtout les rythmes : un crawler peut déverser 500 tâches en trois secondes, les workers les absorberont au rythme qu'autorisent vos threads CaptchaAI.


L'architecture en un coup d'œil

Un seul sens de circulation : les producteurs empilent, les workers dépilent, appellent l'API, écrivent le résultat et préviennent les abonnés.

Producers → Redis List (FIFO queue) → Workers → CaptchaAI API
                                          ↓
                                   Redis Hash (results)
                                          ↓
                                   Redis Pub/Sub (notifications)

Chaque étage est indépendant : ajouter ou redémarrer un worker ne demande aucune reconfiguration ailleurs.


Le gestionnaire de file : soumettre, récupérer, stocker

La classe ci-dessous encapsule toutes les interactions Redis. Une tâche prioritaire part en tête avec lpush, une tâche normale en queue avec rpush : deux niveaux de priorité sans structure supplémentaire. Côté consommation, blpop bloque jusqu'à l'arrivée d'une tâche. Les résultats vont dans un hash doublé d'une clé d'expiration : sans TTL, un pipeline continu sature la mémoire de l'instance.

import json
import time
import uuid
import redis


class CaptchaQueue:
    """Redis-backed CAPTCHA task queue."""

    def __init__(self, redis_url="redis://localhost:6379"):
        self.redis = redis.from_url(redis_url)
        self.queue_key = "captcha:tasks"
        self.results_key = "captcha:results"
        self.notify_channel = "captcha:done"

    def submit(self, method, params, priority="normal"):
        """Submit a CAPTCHA task to the queue."""
        task_id = str(uuid.uuid4())[:8]
        task = {
            "id": task_id,
            "method": method,
            "params": params,
            "submitted_at": time.time(),
            "priority": priority,
        }

        if priority == "high":
            self.redis.lpush(self.queue_key, json.dumps(task))
        else:
            self.redis.rpush(self.queue_key, json.dumps(task))

        return task_id

    def fetch(self, timeout=30):
        """Fetch next task from queue (blocking)."""
        result = self.redis.blpop(self.queue_key, timeout=timeout)
        if result is None:
            return None
        _, raw = result
        return json.loads(raw)

    def store_result(self, task_id, result):
        """Store task result and notify listeners."""
        self.redis.hset(
            self.results_key,
            task_id,
            json.dumps(result),
        )
        # Notify via pub/sub
        self.redis.publish(self.notify_channel, task_id)
        # Set TTL on result (1 hour)
        # Results are in a hash, so we track expiry separately
        self.redis.setex(
            f"captcha:ttl:{task_id}", 3600, "1",
        )

    def get_result(self, task_id):
        """Get result for a task (non-blocking)."""
        raw = self.redis.hget(self.results_key, task_id)
        if raw:
            return json.loads(raw)
        return None

    def wait_result(self, task_id, timeout=120):
        """Wait for a task result via polling."""
        start = time.time()
        while time.time() - start < timeout:
            result = self.get_result(task_id)
            if result:
                return result
            time.sleep(1)
        return None

    def queue_stats(self):
        """Get queue statistics."""
        return {
            "pending": self.redis.llen(self.queue_key),
            "completed": self.redis.hlen(self.results_key),
        }

Le worker : de la tâche Redis au token CaptchaAI

Le worker boucle indéfiniment sur fetch(). Pour chaque tâche, il envoie la requête à in.php, récupère l'identifiant de résolution, puis interroge res.php toutes les 5 secondes tant que la réponse vaut CAPCHA_NOT_READY. N'allez pas plus vite : vous consommeriez du quota de requêtes sans accélérer la résolution.

Les échecs sont stockés comme des résultats, pas propagés en exception : un timeout écrit un statut error dans le hash, le producteur le lit et décide de la nouvelle tentative ; le worker enchaîne.

import os
import time
import requests


class QueueWorker:
    """Worker that processes CAPTCHA tasks from Redis queue."""

    def __init__(self, api_key, queue):
        self.api_key = api_key
        self.queue = queue
        self.base = "https://ocr.captchaai.com"

    def run(self):
        """Main worker loop."""
        worker_id = os.getpid()
        print(f"Worker {worker_id} started")

        while True:
            task = self.queue.fetch(timeout=30)
            if task is None:
                continue

            task_id = task["id"]
            print(f"[{worker_id}] Processing {task_id}")

            start = time.time()
            try:
                token = self._solve(task["method"], task["params"])
                duration = time.time() - start
                self.queue.store_result(task_id, {
                    "status": "success",
                    "token": token,
                    "duration": f"{duration:.1f}s",
                })
                print(f"[{worker_id}] {task_id} solved in {duration:.1f}s")

            except Exception as e:
                self.queue.store_result(task_id, {
                    "status": "error",
                    "error": str(e),
                })
                print(f"[{worker_id}] {task_id} failed: {e}")

    def _solve(self, method, params, timeout=120):
        resp = requests.post(f"{self.base}/in.php", data={
            "key": self.api_key,
            "method": method,
            "json": 1,
            **params,
        }, timeout=30)
        result = resp.json()

        if result.get("status") != 1:
            raise RuntimeError(result.get("request"))

        captcha_id = result["request"]
        start = time.time()

        while time.time() - start < timeout:
            time.sleep(5)
            resp = requests.get(f"{self.base}/res.php", params={
                "key": self.api_key,
                "action": "get",
                "id": captcha_id,
                "json": 1,
            }, timeout=15)
            data = resp.json()
            if data["request"] != "CAPCHA_NOT_READY":
                if data.get("status") == 1:
                    return data["request"]
                raise RuntimeError(data["request"])

        raise TimeoutError("Solve timeout")


# Run worker
if __name__ == "__main__":
    queue = CaptchaQueue()
    worker = QueueWorker(os.environ["CAPTCHAAI_KEY"], queue)
    worker.run()

Lancer plusieurs workers en parallèle

Un worker n'occupe qu'un thread. Le lanceur ci-dessous démarre N processus sur la même file : chacun ouvre sa connexion Redis, et blpop répartit les tâches sans qu'aucune ne soit servie deux fois.

import multiprocessing
import os


def start_workers(num_workers=4):
    """Launch multiple worker processes."""
    queue = CaptchaQueue()
    processes = []

    for i in range(num_workers):
        p = multiprocessing.Process(
            target=run_worker,
            args=(os.environ["CAPTCHAAI_KEY"],),
        )
        p.start()
        processes.append(p)
        print(f"Started worker {i + 1}/{num_workers}")

    return processes


def run_worker(api_key):
    queue = CaptchaQueue()
    worker = QueueWorker(api_key, queue)
    worker.run()


# Launch
processes = start_workers(num_workers=4)

Démarrez avec quatre workers, mesurez, ajustez. Sur une instance Scaleway ou un VPS OVHcloud modeste, plusieurs dizaines de processus tiennent sans peine : ils passent leur temps en attente réseau.


Côté producteur : envoyer un lot et récupérer les tokens

Le producteur n'a besoin que de deux méthodes : submit() pour empiler, wait_result() pour récupérer. Ici, trois pages d'un même parcours utilisateur partent d'un coup, puis les tokens sont collectés dans l'ordre des identifiants.

queue = CaptchaQueue()

# Submit tasks
urls = [
    "https://site1.com/login",
    "https://site2.com/register",
    "https://site3.com/checkout",
]

task_ids = []
for url in urls:
    tid = queue.submit("userrecaptcha", {
        "googlekey": "SITE_KEY",
        "pageurl": url,
    })
    task_ids.append(tid)
    print(f"Submitted {tid} for {url}")

# Wait for all results
for tid in task_ids:
    result = queue.wait_result(tid, timeout=120)
    status = result["status"] if result else "timeout"
    print(f"{tid}: {status}")

# Check queue stats
print(queue.queue_stats())

Le producteur n'attend plus une résolution mais un lot. Sur une recette qui vérifie 40 formulaires, le temps total passe de la somme des résolutions à celui de la plus lente.


Recevoir les résultats en temps réel avec Pub/Sub

Le polling de wait_result() convient à un script ponctuel. Pour un service continu, abonnez-vous au canal de notification : vous êtes prévenu dès que le token est écrit.

import threading


def listen_results(queue):
    """Listen for completed task notifications."""
    pubsub = queue.redis.pubsub()
    pubsub.subscribe(queue.notify_channel)

    for message in pubsub.listen():
        if message["type"] == "message":
            task_id = message["data"].decode()
            result = queue.get_result(task_id)
            print(f"Task {task_id} completed: {result['status']}")


# Run listener in background
listener = threading.Thread(
    target=listen_results,
    args=(CaptchaQueue(),),
    daemon=True,
)
listener.start()

Dimensionner vos workers sur vos threads CaptchaAI

La règle est simple : un worker actif occupe un thread. La facturation CaptchaAI porte sur les threads concurrents, avec un nombre de résolutions illimité par thread ; lancer beaucoup plus de workers que de threads disponibles n'apporte donc rien.

  • BASIC ($15/mois, 5 threads) : pipelines de recette.
  • ADVANCE ($90/mois, 50 threads) : crawl régulier ou suite de tests nocturne.
  • ENTERPRISE ($300/mois, 200 threads) : traitement continu multi-sites.

Exemple : une plateforme e-commerce basée à Lille contrôle chaque nuit les formulaires de connexion et de paiement de ses douze environnements de recette. Avec le plafond de 60 secondes pour reCAPTCHA v2 et 50 threads, 600 tâches se vident en un quart d'heure. Facturation en dollars US.


Exploitation : supervision, reprise et RGPD

Trois réflexes dès la mise en production.

Surveillez la profondeur de file. llen est votre indicateur avancé : une file qui grossit plus de dix minutes signale des producteurs plus rapides que les threads disponibles. Exposez la valeur en métrique et alertez sur un seuil.

Rendez la reprise sûre. blpop retire la tâche avant traitement : si le worker meurt en plein appel, elle disparaît. brpoplpush vers une liste de traitement en cours en conserve une copie, réinjectable au redémarrage.

Minimisez les données journalisées. Les paramètres de tâche contiennent des pageurl, parfois des identifiants de session. Journalisez l'identifiant et le statut, pas le payload, et alignez la rétention Redis sur vos obligations RGPD.


Dépannage

Problème Cause probable Correctif
Workers inactifs, file pleine Connexion sur une autre instance Redis Vérifiez REDIS_URL et la base sélectionnée
Résultats disparus Aucune gestion du TTL Utilisez setex pour l'expiration
File qui ne se vide jamais Workers plus lents que les producteurs Ajoutez des workers ou montez en threads
Tâche traitée deux fois Tâche dépilée puis worker interrompu Passez à brpoplpush
CAPCHA_NOT_READY jusqu'au timeout Polling trop rapproché ou sitekey erroné 5 secondes entre deux appels, revérifiez le sitekey

FAQ

Que se passe-t-il si un worker plante en pleine résolution ?

La tâche est perdue : blpop la retire dès qu'elle est dépilée. Basculez sur brpoplpush, qui la copie dans une liste « en cours » rejouable au démarrage suivant.

Comment empêcher la file de grossir plus vite qu'elle ne se vide ?

Appliquez une limite côté producteur : si llen dépasse 500 tâches, il patiente. Plus fiable que d'ajouter des workers, puisque le plafond réel reste le nombre de threads de votre plan.

Redis managé ou instance auto-hébergée ?

Les deux conviennent : la charge Redis est minuscule face au temps d'attente de l'API. Regardez la proximité réseau — Scaleway, OVHcloud ou une région comme eu-west-3 (Paris).

Faut-il utiliser Redis Streams plutôt que des listes ?

Les Streams apportent les groupes de consommateurs et l'acquittement, ce qui simplifie la reprise sur incident. Les listes restent adaptées à une file simple ; migrez quand plusieurs groupes de lecteurs deviennent nécessaires.

Ce pipeline fonctionne-t-il avec hCaptcha ?

Non — hCaptcha et FunCaptcha ne sont pas pris en charge par CaptchaAI, et GeeTest v4 est annoncé comme à venir. Le même code traite reCAPTCHA v2 et v3, Cloudflare Turnstile, Cloudflare Challenge, GeeTest v3 et les CAPTCHA image : seul method change.


Guides connexes


Répartissez la charge : ouvrez votre compte CaptchaAI et branchez la file.

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