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:
- Der Producer legt eine Aufgabe in
captcha.tasks. - Ein Worker übermittelt sie an CaptchaAI und fragt das Ergebnis ab.
- Das Token geht nach
captcha.results, erst danach bestätigt der Worker. - Abgelehnte Aufgaben wandern über
captcha.dlxincaptcha.failedund 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.