DevOps & Skalierung

RabbitMQ + CaptchaAI: Message Queue-Integration

Ein Worker, der mitten im Lösen abstürzt, darf keine Aufgabe mit sich reißen – darüber entscheidet die Warteschlange, nicht der Solver. RabbitMQ liefert dafür die passenden Bausteine: dauerhafte Queues, eine explizite Bestätigung pro Nachricht und einen Dead-Letter-Exchange für alles, was dauerhaft scheitert. Zusammen mit CaptchaAI entsteht daraus eine CAPTCHA-Pipeline, die Sie neu starten, umziehen und horizontal vergrößern können, ohne anschließend Aufgabenlisten aus Logdateien zu rekonstruieren.

Der Aufbau folgt der Reihenfolge, in der er im Betrieb entsteht: Zustellmodell und Thread-Budget, dann Broker, Producer und Worker, zuletzt Collector, Routing und Monitoring. Alle Beispiele sind in Python mit pika geschrieben und sprechen in.php und res.php unter ocr.captchaai.com an.


Was RabbitMQ im CAPTCHA-Betrieb übernimmt

Der Ablauf ist schlicht:

  1. Der Producer legt eine Aufgabe in captcha.tasks.
  2. Ein Worker übermittelt sie an CaptchaAI und fragt das Ergebnis ab.
  3. Das Token geht nach captcha.results, erst danach bestätigt der Worker.
  4. Abgelehnte Aufgaben wandern über captcha.dlx in captcha.failed und bleiben sichtbar.
RabbitMQ-Baustein Betriebsproblem, das er löst
Dauerhafte Queues (durable=True) Aufgaben überstehen einen Broker-Neustart
Explizite Bestätigung (basic_ack) Ein abgestürzter Worker gibt seine Aufgabe zurück, statt sie zu verlieren
Dead-Letter-Exchange Dauerhaft fehlschlagende Aufgaben landen in captcha.failed statt in einer Endlosschleife
Prioritäten Ein Login-CAPTCHA aus dem laufenden Test geht vor dem Nacht-Batch
Routing-Keys reCAPTCHA, Turnstile und Bild-CAPTCHAs erreichen getrennte Worker-Pools

Thread-Budget planen, bevor der erste Worker startet

CaptchaAI rechnet pro gleichzeitigem Thread ab, nicht pro Lösung – jeder Plan enthält unbegrenzte Lösungen pro Thread. Mit prefetch_count=1 bearbeitet ein Consumer genau eine Aufgabe gleichzeitig. Die Zahl gleichzeitig laufender Consumer sollte Ihre Thread-Zahl daher nicht überschreiten – überzählige Anfragen warten sonst nur auf freie Kapazität.

Plan Preis/Monat Threads Gleichzeitige Consumer (Richtwert)
BASIC 15 $ 5 bis 5
STANDARD 30 $ 15 bis 15
ADVANCE 90 $ 50 bis 50
PREMIUM 170 $ 100 bis 100

Hinweis: Alle Preise in US-Dollar. Maßgeblich ist die Thread-Zahl Ihres Plans, nicht die Zahl der gelösten CAPTCHAs.

Ein Beispiel: Ein Berliner Team zieht nachts Preisdaten aus eigenen Shop-Instanzen und verteilt 30 Worker über drei Hetzner-Cloud-Server. Dazu passt ADVANCE (90 $/Monat, 50 Threads) – mit Luft für Lastspitzen.


RabbitMQ lokal starten

# Docker
docker run -d --hostname rabbitmq \
  -p 5672:5672 -p 15672:15672 \
  rabbitmq:3-management

# Python client
pip install pika requests

Port 5672 bedient das AMQP-Protokoll, Port 15672 die Management-Oberfläche. Der Benutzer guest funktioniert nur über localhost – auf einem Server richten Sie einen eigenen vhost, einen eigenen Benutzer und TLS ein, bevor sich Worker von außen verbinden.


Producer: CAPTCHA-Aufgaben in die Queue legen

import json
import uuid
import pika


class CaptchaProducer:
    """Submit CAPTCHA tasks to RabbitMQ."""

    def __init__(self, rabbitmq_url="amqp://guest:guest@localhost:5672/"):
        self.connection = pika.BlockingConnection(
            pika.URLParameters(rabbitmq_url),
        )
        self.channel = self.connection.channel()
        self._setup_queues()

    def _setup_queues(self):
        """Declare durable queues and exchanges."""
        # Dead letter exchange for failed tasks
        self.channel.exchange_declare(
            exchange="captcha.dlx",
            exchange_type="direct",
            durable=True,
        )
        self.channel.queue_declare(
            queue="captcha.failed",
            durable=True,
        )
        self.channel.queue_bind(
            queue="captcha.failed",
            exchange="captcha.dlx",
            routing_key="failed",
        )

        # Main task queue with dead letter routing
        self.channel.queue_declare(
            queue="captcha.tasks",
            durable=True,
            arguments={
                "x-dead-letter-exchange": "captcha.dlx",
                "x-dead-letter-routing-key": "failed",
                "x-message-ttl": 300000,  # 5 min TTL
            },
        )

        # Results queue
        self.channel.queue_declare(
            queue="captcha.results",
            durable=True,
        )

    def submit(self, method, params, priority=0):
        """Submit a CAPTCHA task."""
        task_id = str(uuid.uuid4())[:8]
        task = {
            "id": task_id,
            "method": method,
            "params": params,
        }

        self.channel.basic_publish(
            exchange="",
            routing_key="captcha.tasks",
            body=json.dumps(task),
            properties=pika.BasicProperties(
                delivery_mode=2,  # Persistent
                priority=priority,
                message_id=task_id,
            ),
        )
        return task_id

    def close(self):
        self.connection.close()


# Usage
producer = CaptchaProducer()

task_id = producer.submit("userrecaptcha", {
    "googlekey": "SITE_KEY",
    "pageurl": "https://example.com",
}, priority=5)

print(f"Submitted: {task_id}")
producer.close()

Drei Details lohnen einen zweiten Blick. delivery_mode=2 schreibt die Nachricht auf die Platte – sonst nützt die dauerhafte Queue wenig. Das x-message-ttl von 300.000 ms verwirft Aufgaben, die länger als fünf Minuten liegen bleiben; ein zu spät geliefertes Token ist ohnehin wertlos, weil es nur rund 120 s gültig ist. Und priority greift erst, wenn die Queue mit x-max-priority deklariert wurde: Sonst nimmt RabbitMQ die Eigenschaft zwar an, sortiert aber weiter nach Eingangsreihenfolge.


Consumer: der Worker, der die CAPTCHAs löst

import json
import os
import time
import pika
import requests


class CaptchaConsumer:
    """RabbitMQ consumer that solves CAPTCHAs."""

    def __init__(self, api_key, rabbitmq_url="amqp://guest:guest@localhost:5672/"):
        self.api_key = api_key
        self.base = "https://ocr.captchaai.com"
        self.connection = pika.BlockingConnection(
            pika.URLParameters(rabbitmq_url),
        )
        self.channel = self.connection.channel()
        # Process one task at a time
        self.channel.basic_qos(prefetch_count=1)

    def start(self):
        """Start consuming tasks."""
        self.channel.basic_consume(
            queue="captcha.tasks",
            on_message_callback=self._handle_task,
        )
        print("Worker started. Waiting for tasks...")
        self.channel.start_consuming()

    def _handle_task(self, ch, method, properties, body):
        """Process a single CAPTCHA task."""
        task = json.loads(body)
        task_id = task["id"]
        print(f"Processing {task_id}...")

        try:
            token = self._solve(task["method"], task["params"])

            # Publish result
            result = {
                "task_id": task_id,
                "status": "success",
                "token": token,
            }
            ch.basic_publish(
                exchange="",
                routing_key="captcha.results",
                body=json.dumps(result),
                properties=pika.BasicProperties(delivery_mode=2),
            )

            # Acknowledge message (remove from queue)
            ch.basic_ack(delivery_tag=method.delivery_tag)
            print(f"{task_id} solved successfully")

        except Exception as e:
            print(f"{task_id} failed: {e}")

            # Reject and send to dead letter queue
            ch.basic_nack(
                delivery_tag=method.delivery_tag,
                requeue=False,  # Goes to DLX
            )

    def _solve(self, captcha_method, params, timeout=120):
        resp = requests.post(f"{self.base}/in.php", data={
            "key": self.api_key,
            "method": captcha_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__":
    consumer = CaptchaConsumer(os.environ["CAPTCHAAI_KEY"])
    consumer.start()

Die Reihenfolge im Callback ist der Kern der Ausfallsicherheit: erst das Token nach captcha.results schreiben, dann bestätigen. Bricht der Prozess vorher ab, stellt RabbitMQ die unbestätigte Nachricht dem nächsten Consumer zu. Ein basic_nack mit requeue=False schickt die Aufgabe dagegen sofort in die Dead-Letter-Queue – richtig so, denn ein falscher sitekey wird beim zehnten Versuch nicht plötzlich korrekt. Der API-Schlüssel kommt aus der Umgebungsvariablen CAPTCHAAI_KEY und gehört weder ins Image noch ins Repository.


Ergebnisse aus der Results-Queue einsammeln

import json
import pika


class ResultCollector:
    """Collect task results from the results queue."""

    def __init__(self, rabbitmq_url="amqp://guest:guest@localhost:5672/"):
        self.connection = pika.BlockingConnection(
            pika.URLParameters(rabbitmq_url),
        )
        self.channel = self.connection.channel()
        self.results = {}

    def collect(self, expected_count, timeout=120):
        """Collect a specific number of results."""
        deadline = time.time() + timeout

        while len(self.results) < expected_count and time.time() < deadline:
            method, _, body = self.channel.basic_get(
                queue="captcha.results",
                auto_ack=True,
            )
            if body:
                result = json.loads(body)
                self.results[result["task_id"]] = result

            time.sleep(0.5)

        return self.results

Das Beispiel nutzt time, ohne es zu importieren – ergänzen Sie den Import. basic_get fragt aktiv ab und passt zu abgeschlossenen Batches, etwa einem CI-Job, der auf eine feste Zahl Tokens wartet. Im Dauerbetrieb ist basic_consume besser: Der Broker stellt zu, statt dass Ihr Prozess alle 0,5 Sekunden nachfragt.


Routing nach CAPTCHA-Typ

CaptchaAI löst unter anderem reCAPTCHA v2 und v3, Cloudflare Turnstile und Cloudflare Challenge, GeeTest v3 sowie Bild- und Rasterbild-CAPTCHAs; hCaptcha und FunCaptcha gehören nicht dazu, GeeTest v4 ist als bald verfügbar angekündigt. CaptchaFox, Friendly Captcha und Lemin laufen im Beta-Status.

# Setup exchanges and queues
channel.exchange_declare(
    exchange="captcha.types",
    exchange_type="direct",
    durable=True,
)

# Queue per type
for captcha_type in ["recaptcha", "turnstile", "image"]:
    channel.queue_declare(queue=f"captcha.{captcha_type}", durable=True)
    channel.queue_bind(
        queue=f"captcha.{captcha_type}",
        exchange="captcha.types",
        routing_key=captcha_type,
    )


# Submit with routing
def submit_routed(channel, captcha_type, task):
    channel.basic_publish(
        exchange="captcha.types",
        routing_key=captcha_type,
        body=json.dumps(task),
        properties=pika.BasicProperties(delivery_mode=2),
    )

Getrennte Queues lohnen sich wegen der Lösungszeit: Bild-CAPTCHAs sind in unter 0,5 s durch, Cloudflare Turnstile in unter 10 s, reCAPTCHA v2 kann bis zu 60 s brauchen. In einer gemeinsamen Queue bremst der langsame Typ die schnellen aus, und getrennte Pools skalieren Sie unabhängig – zwei Worker für Bilder, acht für reCAPTCHA v2.


Betrieb: Monitoring, Deployment und Datenschutz

Im Livebetrieb genügen vier Kennzahlen aus der Management-Oberfläche oder dem Prometheus-Plugin:

  • Länge von captcha.tasks – steigt sie dauerhaft, fehlen Consumer oder Threads.
  • Unbestätigte Nachrichten – bleiben sie konstant hoch, hängt ein Worker in einem Timeout.
  • Zulauf auf captcha.failed – jede Nachricht dort ist ein Parameterfehler oder ein API-Fehlercode, den jemand ansehen sollte.
  • Verbindungsabbrüche – meist ein zu kurzes Heartbeat-Intervall bei langen Lösungszeiten.

Der Worker selbst ist zustandslos: eine systemd-Unit oder ein Container pro Prozess, Skalierung über die Replikatzahl, Rollout über GitLab CI. Server bei Hetzner, IONOS oder netcup eignen sich dafür genauso wie eine bestehende Kubernetes-Umgebung.

Ein Punkt, der in DACH-Projekten regelmäßig zu Rückfragen führt: Die Payload liegt im Broker und in dessen Sicherungen. Schreiben Sie nur die nötigen Felder in die Queue und klären Sie die Rechtsgrundlage, wenn personenbezogene Daten wie IP-Adressen darin vorkommen.


Typische Fehlerbilder und ihre Ursachen

Symptom Ursache Abhilfe
Aufgaben fehlen nach einem Neustart Queue oder Nachricht nicht dauerhaft durable=True und delivery_mode=2 setzen
Ein Worker hängt an einer einzigen Aufgabe Lange Lösungszeit, zu hoher Prefetch prefetch_count=1 pro Consumer
captcha.failed wächst stetig Wiederkehrende Parameter- oder Schlüsselfehler Nachrichten dort auswerten und Ursache beheben
Verbindung bricht regelmäßig ab Heartbeat-Timeout Heartbeat-Intervall setzen und Wiederverbindungslogik ergänzen
Prioritäten wirken nicht x-max-priority fehlt bei der Deklaration Queue mit Prioritätsargument neu anlegen

FAQ

Was passiert mit einer Aufgabe, wenn der Worker mitten im Lösen abstürzt?

Sie geht zurück in die Queue. Solange kein basic_ack gesendet wurde, gilt die Nachricht als unbestätigt; bricht die Verbindung ab, stellt RabbitMQ sie dem nächsten freien Consumer zu.

Wie viele Consumer passen zu meinem CaptchaAI-Plan?

So viele, wie der Plan Threads enthält. Mit prefetch_count=1 entspricht ein laufender Consumer einem gleichzeitigen Solve – bei STANDARD (30 $/Monat, 15 Threads) sind das bis zu 15 Prozesse.

Warum werden meine prioritären Aufgaben nicht zuerst gelöst?

Weil die Queue ohne x-max-priority keine Prioritäten kennt. Legen Sie sie mit diesem Argument neu an; eine bestehende Queue lässt sich nachträglich nicht umdeklarieren.

Wie lange ist ein gelöstes Token gültig?

Rund 120 Sekunden. Deshalb sollte zwischen Lösung und Absenden des Formulars möglichst wenig Zeit liegen – lieber knapp vor der Übermittlung lösen als Tokens auf Vorrat sammeln.

Läuft dasselbe Setup auch in Kubernetes?

Ja. Der Consumer hält keinen lokalen Zustand; ein Deployment mit passender Replikatzahl genügt. Achten Sie auf sauberes Beenden, damit ein Pod seine Aufgabe noch bestätigen oder zurückgeben kann.


Verwandte Leitfäden


Queue-gestützte CAPTCHA-Verarbeitung mit RabbitMQ und CaptchaAI.

Kommentare sind für diesen Artikel deaktiviert.