DevOps & Skalierung

NATS-Messaging + CaptchaAI: Leichte CAPTCHA-Aufgabenverteilung

Um CAPTCHA-Aufgaben auf mehrere Worker zu verteilen, brauchen Sie keinen Broker-Cluster mit ZooKeeper und eigenem Betriebsteam. Ein einzelner nats-server-Prozess mit rund 20 MB Speicherbedarf reicht aus: Ihre Scraper publizieren Aufgaben auf ein Subject, eine Queue-Group verteilt sie auf beliebig viele Worker, und jeder Worker übergibt die Abfrage an CaptchaAI und schreibt das Token zurück.

Das funktioniert, weil das Aufgabenprofil dazu passt: CAPTCHA-Aufgaben sind kurzlebig. Ein Token ist nur wenige Minuten gültig, eine verlorene Aufgabe kein Datenverlust, sondern ein erneuter Versuch. Sie bezahlen also keinen Betriebsaufwand für Zustellgarantien, die Sie gar nicht brauchen. Dieser Leitfaden zeigt den vollständigen Weg – Architektur, Publisher, Worker-Pool, Ergebnis-Collector, optionale Persistenz über JetStream – mit lauffähigem Code in Python und Node.js.

Architektur: drei Subjects, ein Worker-Pool

[Scrapers] → Publish → [NATS: captcha.tasks]
                              ↓
                    Queue Group: captcha-workers
                    ├── Worker 1 (solve via CaptchaAI)
                    ├── Worker 2
                    └── Worker 3
                              ↓
                    Publish → [NATS: captcha.results]
                              ↓
                    [Result Subscribers]

Das entscheidende Detail ist die Queue-Group captcha-workers: NATS liefert jede Nachricht an genau ein Mitglied der Gruppe. Sie starten zusätzliche Worker, ohne Partitionen neu zu verteilen – wer sich anmeldet, bekommt Arbeit; wer ausfällt, wird übergangen.

Zwei Subjects genügen für den Regelbetrieb: captcha.tasks für eingehende Aufgaben und captcha.results für gelöste Tokens und Fehler. Ein drittes Subject für endgültig fehlgeschlagene Aufgaben lohnt sich erst, wenn Sie Wiederholungen protokollieren müssen.

NATS, Kafka oder RabbitMQ – die Entscheidung in einer Tabelle

Merkmal NATS Kafka RabbitMQ
Latenz < 1 ms 5–10 ms 1–5 ms
Aufwand beim Setup einzelne Binärdatei Cluster + ZooKeeper mittel
Speicherbedarf ~20 MB ~1 GB+ ~200 MB
Persistenz optional (JetStream) eingebaut eingebaut
Am besten geeignet für kurzlebige Aufgaben, niedrige Latenz dauerhaftes Streaming komplexes Routing

Hinweis: Die Werte sind Größenordnungen aus üblichen Standardkonfigurationen, keine Benchmarks – messen Sie mit eigenem Traffic-Muster nach.

Als Faustregel:

  • NATS, wenn minimale Infrastruktur und niedrige Latenz zählen.
  • Kafka, wenn Aufgaben dauerhaft gespeichert und erneut abgespielt werden müssen.
  • RabbitMQ, wenn Sie Routing mit Prioritäten und Exchanges brauchen.

Wie viele Worker Ihr Thread-Kontingent verträgt

CaptchaAI rechnet pro gleichzeitigem Thread ab, nicht pro Lösung – innerhalb eines Plans sind die Lösungen pro Thread unbegrenzt. Für die Dimensionierung des Worker-Pools heißt das: Orientieren Sie sich am Thread-Kontingent, nicht an der Queue-Tiefe. Mit STANDARD (30 $/Monat, 15 Threads) laufen 15 Abfragen parallel; ein 30. Worker bringt dann keinen Durchsatz mehr, sondern nur Wartezeit. Wer dauerhaft mehr Parallelität braucht, wechselt zu ADVANCE (90 $/Monat, 50 Threads) oder PREMIUM (170 $/Monat, 100 Threads). Preise in US-Dollar.

Die zweite Größe ist die Lösungszeit pro CAPTCHA-Typ: Cloudflare Turnstile ist üblicherweise in unter 10 Sekunden erledigt, reCAPTCHA v2 kann bis zu 60 Sekunden dauern. Ein Pool für reCAPTCHA v2 hält seine Threads bei gleicher Aufgabenzahl also deutlich länger belegt als einer für Turnstile.

Voraussetzungen: Server und Clients installieren

Der Server ist eine einzelne Binärdatei ohne externe Abhängigkeiten und startet mit der Standardkonfiguration auf Port 4222.

# 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

Schritt 1: Aufgaben publizieren (Scraper-Seite)

Der Publisher bleibt bewusst simpel: Er serialisiert die Aufgabe als JSON und schickt sie an captcha.tasks. Alles, was der Worker später braucht, muss in der Payload stehen – Methode, Sitekey und Page-URL. Der Aufruf await nc.flush() am Ende ist nicht optional: Ohne ihn verlieren kurzlebige Skripte die zuletzt gepufferten Nachrichten beim Verbindungsabbau.

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

Schritt 2: Worker als Queue-Group-Abonnent

Der Worker erledigt drei Dinge:

  1. Die Aufgabe aus der Queue-Group entgegennehmen.
  2. Sie an in.php übermitteln.
  3. Das Ergebnis über res.php abfragen.

Solange die API CAPCHA_NOT_READY meldet, fragt der Worker den Status erneut ab; jede andere Antwort beendet den Versuch – entweder mit Token oder mit Fehlercode. Der API-Schlüssel kommt aus der Umgebungsvariablen CAPTCHAAI_API_KEY und gehört nicht in den Quelltext, schon gar nicht in ein Container-Image.

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

Schritt 3: Ergebnisse einsammeln

Der Collector abonniert captcha.results ohne Queue-Group – jeder Subscriber sieht damit jedes Ergebnis. Das ist praktisch, wenn ein Prozess die Tokens zurück in die Scraper-Pipeline gibt und ein zweiter nur Kennzahlen mitschreibt.

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

JetStream: wenn Aufgaben nicht verloren gehen dürfen

Core-NATS puffert nicht für langsame Consumer. Wenn Ihre Aufgaben eine verlässliche Zustellung brauchen, aktivieren Sie JetStream – die Persistenzschicht des Servers, für die Sie weder zusätzliche Dienste noch ein zweites Betriebskonzept einführen müssen.

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 ergänzt Festplattenpersistenz, Replay und bestätigte Zustellung – funktional nahe an Kafka, im Betrieb aber schlanker. Der Aufwand bleibt: Streams wollen dimensioniert, Retention-Regeln und Ack-Timeouts gesetzt werden. Für reine CAPTCHA-Aufgaben lohnt sich das selten, für Abläufe mit Nachweispflicht schon.

Worker horizontal skalieren

# 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

Skalieren Sie in kleinen Schritten und beobachten Sie dabei zwei Kennzahlen: die Anzahl offener Aufgaben und die Lösungszeit pro Worker. Steigt die Lösungszeit, während die Warteschlange gleich lang bleibt, ist nicht NATS der Engpass, sondern Ihr Thread-Kontingent.

Aus der Praxis: verteilte Worker auf DACH-Infrastruktur

Ein typischer Aufbau aus dem deutschsprachigen Raum sieht so aus: Ein Hamburger Team für Preis-Monitoring betreibt den nats-server auf einer kleinen Hetzner-VM in Nürnberg und fährt drei Worker als Container auf demselben Host; ausgerollt wird über eine GitLab-CI-Pipeline, die ohnehin für den Rest des Stacks läuft.

Zwei Punkte tauchen dabei regelmäßig auf:

  • Netzwerklatenz: Gegenüber Lösungszeiten von mehreren Sekunden fällt der Weg zur API kaum ins Gewicht – ein Standortwechsel bringt wenig.
  • Datenschutz: Logdaten mit IP-Adressen und vollständigen Page-URLs gelten unter der DSGVO regelmäßig als personenbezogene Daten. Prüfen Sie, was Sie auf captcha.results mitschreiben und wie lange Sie es aufbewahren.

Typische Fehlerbilder

Symptom Ursache Vorgehen
Nachrichten verschwinden unter Last Core-NATS puffert nicht für langsame Consumer JetStream aktivieren oder Worker-Zahl erhöhen
Worker läuft, empfängt aber nichts Subject oder Name der Queue-Group weicht vom Publisher ab beide Strings vergleichen: captcha.tasks und captcha-workers
Verbindung bricht regelmäßig ab Server-Neustart oder Netzwerkunterbrechung Auto-Reconnect in den Client-Optionen aktivieren
Ein Worker bekommt deutlich mehr Aufgaben schnellere Worker holen sich mehr Nachrichten normales Verhalten; eingreifen erst, wenn ein Worker gar nichts erhält
Ergebnisse enthalten Fehler statt Tokens falscher API-Schlüssel oder aufgebrauchtes Guthaben Schlüssel und Guthaben im Dashboard prüfen, dann Fehlercodes auswerten

Häufige Fragen

Wie viele Worker sollte ich starten?

So viele, wie Ihr Thread-Kontingent hergibt – zusätzliche Prozesse warten sonst nur auf freie Threads. Starten Sie mit einem Worker pro fünf Threads und erhöhen Sie schrittweise, solange die Lösungszeit stabil bleibt.

Was passiert, wenn ein Worker mitten in der Verarbeitung abstürzt?

Mit Core-NATS ist die Aufgabe verloren, weil die Nachricht bereits zugestellt wurde – der Publisher muss sie erneut senden. Mit JetStream und einem Ack-Timeout wird die Nachricht automatisch an einen anderen Worker ausgeliefert.

Lässt sich der Aufbau in Docker oder Kubernetes betreiben?

Ja. Der Server läuft als schlanker Container, die Worker als eigenes Deployment; skaliert wird über die Replica-Zahl. Die Queue-Group erledigt die Lastverteilung, ein zusätzlicher Load Balancer ist nicht nötig.

Wie lange bleibt ein gelöstes Token brauchbar?

Nur kurz – bei reCAPTCHA rund 120 Sekunden. Lösen Sie die Abfrage deshalb unmittelbar vor dem Absenden des Formulars, statt Tokens zwischenzuspeichern – eine lange Verweildauer in der Warteschlange produziert sonst abgelaufene Tokens.


Verwandte Leitfäden

Kommentare sind für diesen Artikel deaktiviert.