DevOps & Mise à l'Échelle

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

Redis fournit une mise en file d'attente rapide et fiable pour distribuer les tâches CAPTCHA entre plusieurs travailleurs. Ce guide construit un système producteur-consommateur complet avec suivi des résultats et gestion des erreurs.


Architecture

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

Gestionnaire de file d'attente de tâches

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),
        }

Processus de travail

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

Lanceur multi-travailleurs

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)

Exemple de producteur

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

Écouteur de résultats Pub/Sub

Soyez averti immédiatement lorsque les tâches sont terminées :

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

Dépannage

Problème Parce que Corriger
Les travailleurs sont inactifs avec des tâches en file d'attente Connexion à un mauvais Redis Vérifiez REDIS_URL
Les résultats disparaissent Pas de gestion TTL Utilisez setex pour l'expiration des résultats
La file d'attente s'étend sans limite Les ouvriers sont trop lents Ajoutez plus de travailleurs ou augmentez la simultanéité
Traitement en double La tâche a sauté mais le travailleur plante Utilisez brpoplpush pour une file d'attente fiable

FAQ

Pourquoi Redis plutôt que d'autres files d'attente de messages ?

Redis est simple, rapide et la plupart des équipes l'exécutent déjà. Pour un routage complexe ou une livraison garantie, pensez plutôt à RabbitMQ.

Combien de travailleurs par instance Redis ?

Une seule instance Redis gère plus de 100 000 opérations/second. Le goulot d'étranglement est le débit de l'API CaptchaAI, et non Redis. Planifiez les travailleurs en fonction de votre capacité CaptchaAI.

Dois-je utiliser Redis Streams au lieu de listes ?

Redis Streams fournit des groupes de consommateurs et une reconnaissance, ce qui est meilleur pour la production. Les listes fonctionnent bien pour les configurations plus simples.


Guides connexes

  • LapinMQ + CaptchaAI
  • Résolution de CAPTCHA par lots

Répartissez votre charge de travail -obtenir CaptchaAIaujourd'hui.

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