DevOps & Skalierung

Apache Kafka + CaptchaAI: Streaming-CAPTCHA-Aufgabenverarbeitung

Apache Kafka lohnt sich für die CAPTCHA-Verarbeitung genau dann, wenn Ihr Durchsatz die Grenzen einer einfachen Warteschlange sprengt – also ab mehreren Tausend Aufgaben pro Stunde, verteilt auf mehrere Worker. Kafka entkoppelt die Übermittlung von CAPTCHA-Aufgaben sauber von der Ergebnisverarbeitung: Producer schreiben Aufgaben in ein Topic, eine Consumer-Group löst sie über CaptchaAI, und die Antworten landen in einem zweiten Topic für die nachgelagerte Verarbeitung.

Unterhalb dieser Größenordnung ist Kafka meist überdimensioniert. Wer pro Tag nur wenige Hundert CAPTCHAs löst, fährt mit einer schlankeren Queue günstiger und betriebsärmer. Interessant wird Kafka erst, wenn Sie Rebalancing, Consumer-Groups, Replays und klar getrennte Producer- und Consumer-Pfade tatsächlich benötigen – etwa bei einem verteilten Scraping-Cluster, dessen Last über den Tag stark schwankt.

Kafka oder einfache Queue?

Betriebsrealität Kafka passt gut Einfachere Queue reicht oft
Tausende Aufgaben pro Stunde, mehrere Konsumenten Ja
Replay, Audit oder Streaming-Auswertung sind wichtig Ja
Ein kleiner interner Workflow mit begrenztem Volumen Häufig ausreichend
Sie brauchen nur eine robuste Aufgabenliste, keine Streaming-Plattform Häufig ausreichend

Als Faustregel: Erst wenn eine Queue allein Ihre Anforderungen an Nachvollziehbarkeit und horizontale Skalierung nicht mehr erfüllt, spielt Kafka seine Stärken aus. Konkrete Signale, die für Kafka sprechen:

  • Sie müssen dieselben Aufgaben zu Analyse- oder Audit-Zwecken erneut abspielen (Replay).
  • Mehrere unabhängige Teams oder Dienste konsumieren denselben Ergebnisstrom.
  • Ihr Volumen schwankt stark und Sie wollen Worker elastisch nachziehen.

Architektur im Überblick

[Scrapers] → Produce → [Kafka: captcha-tasks topic]
                              ↓
                    [CAPTCHA Worker Group]
                    (consume tasks, solve via CaptchaAI)
                              ↓
                    Produce → [Kafka: captcha-results topic]
                              ↓
                    [Result Consumer Group]
                    (process solutions, update database)

Zwei Topics trennen die Zuständigkeiten sauber voneinander:

  • captcha-tasks – CAPTCHA-Parameter, die auf ihre Lösung warten
  • captcha-results – gelöste Token, bereit für die nachgelagerte Verwendung

Der Producer muss nichts über die Ergebnis-Consumer wissen – und umgekehrt.

Hinweis: Genau diese Entkopplung macht die Pipeline unabhängig skalierbar. Sie können die Scraper-Seite und die Solver-Seite getrennt deployen, versionieren und hochskalieren, ohne dass eine Änderung die andere blockiert.

Voraussetzungen

# Python
pip install kafka-python requests

# Node.js
npm install kafkajs axios

Dazu ein Kafka-Broker, der auf localhost:9092 läuft (oder unter Ihrer Cluster-Adresse). In der Praxis laufen Broker und Worker in der DACH-Region häufig auf Infrastruktur wie Hetzner oder netcup.

DSGVO-Hinweis: Bei Scraping-Workloads gelten IP-Adressen als personenbezogene Daten. Prüfen Sie Rechtsgrundlage und Datenflüsse, bevor Sie in den Produktivbetrieb gehen – das ist Sorgfaltspflicht auf Ihrer Seite, nicht Teil der CaptchaAI-Integration.

Schritt 1: Topics anlegen

kafka-topics.sh --create --topic captcha-tasks \
  --partitions 6 --replication-factor 1 \
  --bootstrap-server localhost:9092

kafka-topics.sh --create --topic captcha-results \
  --partitions 6 --replication-factor 1 \
  --bootstrap-server localhost:9092

Sechs Partitionen erlauben bis zu sechs parallele Consumer pro Group. Die Partitionszahl ist damit Ihre Obergrenze für echte Parallelität – planen Sie sie an Ihrer Spitzenlast aus.

Tipp: Wählen Sie die Partitionszahl von Anfang an großzügig. Sie lässt sich zwar später erhöhen, doch das verändert die Zuordnung von Schlüssel zu Partition und damit die Reihenfolgegarantien innerhalb eines Schlüssels.

Schritt 2: Task-Producer (Scraper-Seite)

Der Producer schreibt jede CAPTCHA-Aufgabe mit einem Schlüssel in das Topic. Der Schlüssel sorgt dafür, dass zusammengehörige Aufgaben in derselben Partition landen.

Zur Konfiguration: acks="all" sorgt dafür, dass eine Aufgabe erst als bestätigt gilt, wenn alle Replikate sie geschrieben haben – der sicherste Modus gegen Datenverlust. In Kombination mit retries fängt der Producer kurze Broker-Aussetzer selbst ab.

Python

import json
from kafka import KafkaProducer

producer = KafkaProducer(
    bootstrap_servers=["localhost:9092"],
    value_serializer=lambda v: json.dumps(v).encode("utf-8"),
    key_serializer=lambda k: k.encode("utf-8") if k else None,
    acks="all",  # Wait for all replicas to confirm
    retries=3
)


def enqueue_captcha(task_id, sitekey, pageurl, captcha_type="userrecaptcha"):
    """Send a CAPTCHA task to Kafka."""
    task = {
        "task_id": task_id,
        "method": captcha_type,
        "sitekey": sitekey,
        "pageurl": pageurl,
        "submitted_at": __import__("time").time()
    }

    future = producer.send(
        "captcha-tasks",
        key=task_id,  # Key ensures same task goes to same partition
        value=task
    )
    future.get(timeout=10)  # Block until confirmed
    return task_id


# Submit tasks
enqueue_captcha("task_001", "6Le-wvkSAAAAAPBMRTvw0Q4Muexq9bi0DJwx_mJ-", "https://example.com")
enqueue_captcha("task_002", "6Le-wvkSAAAAAPBMRTvw0Q4Muexq9bi0DJwx_mJ-", "https://example.com")
producer.flush()

JavaScript

const { Kafka } = require("kafkajs");

const kafka = new Kafka({
  clientId: "captcha-producer",
  brokers: ["localhost:9092"],
});

const producer = kafka.producer();

async function enqueueCaptcha(taskId, sitekey, pageurl) {
  await producer.connect();

  const task = {
    task_id: taskId,
    method: "userrecaptcha",
    sitekey: sitekey,
    pageurl: pageurl,
    submitted_at: Date.now(),
  };

  await producer.send({
    topic: "captcha-tasks",
    messages: [{ key: taskId, value: JSON.stringify(task) }],
  });
}

(async () => {
  await enqueueCaptcha(
    "task_001",
    "6Le-wvkSAAAAAPBMRTvw0Q4Muexq9bi0DJwx_mJ-",
    "https://example.com"
  );
  await producer.disconnect();
})();

Schritt 3: CAPTCHA-Worker (Consumer + Solver)

Der Worker liest Aufgaben aus captcha-tasks, übermittelt sie an CaptchaAI, fragt das Ergebnis ab und schreibt es nach captcha-results.

Wichtig: Der Offset wird bewusst erst nach erfolgreicher Verarbeitung committet (enable_auto_commit=False plus manueller commit()). So geht bei einem Absturz keine Aufgabe verloren – sie wird von einem anderen Worker erneut gelesen.

Python

import json
import os
import time
import requests
from kafka import KafkaConsumer, KafkaProducer

API_KEY = os.environ["CAPTCHAAI_API_KEY"]

consumer = KafkaConsumer(
    "captcha-tasks",
    bootstrap_servers=["localhost:9092"],
    group_id="captcha-workers",
    value_deserializer=lambda m: json.loads(m.decode("utf-8")),
    auto_offset_reset="earliest",
    enable_auto_commit=False,  # Manual commit after processing
    max_poll_records=10
)

result_producer = KafkaProducer(
    bootstrap_servers=["localhost:9092"],
    value_serializer=lambda v: json.dumps(v).encode("utf-8")
)


def solve_captcha(task):
    """Submit to CaptchaAI and poll for result."""
    # Submit
    resp = requests.post("https://ocr.captchaai.com/in.php", data={
        "key": API_KEY,
        "method": task["method"],
        "googlekey": task["sitekey"],
        "pageurl": task["pageurl"],
        "json": 1
    })
    data = resp.json()

    if data.get("status") != 1:
        return {"error": data.get("request")}

    captcha_id = data["request"]

    # Poll for result
    for _ in range(60):
        time.sleep(5)
        result = requests.get("https://ocr.captchaai.com/res.php", params={
            "key": API_KEY,
            "action": "get",
            "id": captcha_id,
            "json": 1
        }).json()

        if result.get("status") == 1:
            return {"solution": result["request"]}
        if result.get("request") != "CAPCHA_NOT_READY":
            return {"error": result.get("request")}

    return {"error": "TIMEOUT"}


# Main consumer loop
print("CAPTCHA worker started. Waiting for tasks...")
for message in consumer:
    task = message.value
    print(f"Processing {task['task_id']}...")

    result = solve_captcha(task)
    result["task_id"] = task["task_id"]
    result["solved_at"] = time.time()

    # Publish result
    result_producer.send("captcha-results", value=result)
    result_producer.flush()

    # Commit offset after successful processing
    consumer.commit()
    print(f"  → {task['task_id']}: {'solved' if 'solution' in result else result.get('error')}")

JavaScript

const { Kafka } = require("kafkajs");
const axios = require("axios");

const API_KEY = process.env.CAPTCHAAI_API_KEY;

const kafka = new Kafka({
  clientId: "captcha-worker",
  brokers: ["localhost:9092"],
});

const consumer = kafka.consumer({ groupId: "captcha-workers" });
const producer = kafka.producer();

function sleep(ms) {
  return new Promise((resolve) => setTimeout(resolve, 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 { 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 { solution: result.data.request };
    if (result.data.request !== "CAPCHA_NOT_READY")
      return { error: result.data.request };
  }

  return { error: "TIMEOUT" };
}

async function run() {
  await consumer.connect();
  await producer.connect();
  await consumer.subscribe({ topic: "captcha-tasks", fromBeginning: false });

  await consumer.run({
    eachMessage: async ({ message }) => {
      const task = JSON.parse(message.value.toString());
      console.log(`Processing ${task.task_id}...`);

      const result = await solveCaptcha(task);
      result.task_id = task.task_id;
      result.solved_at = Date.now();

      await producer.send({
        topic: "captcha-results",
        messages: [{ value: JSON.stringify(result) }],
      });

      console.log(
        `  → ${task.task_id}: ${result.solution ? "solved" : result.error}`
      );
    },
  });
}

run();

Worker skalieren

Consumer-Groups verteilen die Partitionen automatisch auf alle laufenden Worker:

# 6 partitions, 3 workers → each worker gets 2 partitions
Worker-1: partitions 0, 1
Worker-2: partitions 2, 3
Worker-3: partitions 4, 5

# Add Worker-4 → rebalance
Worker-1: partitions 0, 1
Worker-2: partitions 2
Worker-3: partitions 3, 4
Worker-4: partition 5

Sie können bis zur Anzahl der Partitionen hochskalieren – ein siebter Worker bei sechs Partitionen bleibt ohne Zuteilung leer. Brauchen Sie mehr Parallelität, erhöhen Sie zuerst die Partitionszahl.

Wichtig ist die Abstimmung mit Ihrer CaptchaAI-Thread-Zuteilung: CaptchaAI rechnet pro gleichzeitigem Thread ab, nicht pro Lösung.

Threads statt Lösungen: Ein Plan wie PREMIUM (170 $/Monat, 100 Threads) erlaubt bis zu 100 gleichzeitig laufende Lösungen – mit unbegrenzten Lösungen pro Thread im Abrechnungsmonat. Mehr Kafka-Worker als verfügbare Threads bringen keinen zusätzlichen Durchsatz, sondern nur Wartezeit an der API. Skalieren Sie Worker und Thread-Kontingent deshalb gemeinsam.

Monitoring: Consumer-Lag im Blick

Der wichtigste Frühindikator ist der Consumer-Lag – der Rückstand zwischen produzierten und verarbeiteten Nachrichten:

kafka-consumer-groups.sh --describe --group captcha-workers \
  --bootstrap-server localhost:9092
Metrik Gesund Warnung
Consumer-Lag < 100 > 1000 (Worker hinzufügen)
Nachrichten/s eingehend Entspricht der Scraping-Rate Spitzen deuten auf Bursts hin
Nachrichten/s ausgehend Entspricht der Eingangsrate Rückstand = Engpass

Ein dauerhaft wachsender Lag ist das klarste Signal, dass Ihre Worker mit der Aufgabenrate nicht mehr Schritt halten.

Alarmierung: Setzen Sie einen Schwellenwert auf den Consumer-Lag statt auf die reine Nachrichtenrate. Ein kurzer Burst ist unkritisch, solange der Lag danach wieder abgebaut wird; erst ein anhaltend hoher Lag rechtfertigt zusätzliche Worker oder Threads.

Fehlerbehebung

Die meisten Probleme im Produktivbetrieb lassen sich auf drei Ursachen zurückführen: zu wenige Worker, unsauberes Offset-Handling oder instabile Consumer, die ständiges Rebalancing auslösen. Die folgende Tabelle ordnet Symptome den passenden Gegenmaßnahmen zu.

Problem Ursache Lösung
Consumer-Lag wächst stetig Worker kommen mit der Aufgabenrate nicht nach Weitere Worker-Instanzen ergänzen (bis zur Partitionszahl), danach Partitionen erhöhen
Doppelte Ergebnisse Worker stürzt vor dem Offset-Commit ab Idempotenzprüfung auf task_id im Ergebnis-Consumer
Rebalancing zu häufig Worker stürzen ab oder starten ständig neu session.timeout.ms erhöhen; auf OOM prüfen
Aufgaben ungleich verteilt Schlechte Schlüsselverteilung Zufällige Schlüssel verwenden oder Partitionszahl erhöhen

Häufige Fragen

Die folgenden Fragen tauchen beim Betrieb einer Kafka-gestützten CAPTCHA-Pipeline am häufigsten auf.

Wie viele Partitionen sollte das captcha-tasks-Topic haben?

So viele, wie Sie maximal parallele Worker betreiben wollen. Die Partitionszahl ist die Obergrenze für echte Parallelität in einer Consumer-Group – sechs Partitionen bedeuten höchstens sechs gleichzeitig arbeitende Worker. Planen Sie sie an Ihrer Spitzenlast aus, denn nachträgliches Erhöhen ist möglich, aber verändert die Schlüsselverteilung.

Was passiert mit einem Token, das lange in Kafka liegt?

reCAPTCHA- und Turnstile-Token laufen typischerweise nach etwa 120 Sekunden ab. Lösen Sie deshalb möglichst nah am Zeitpunkt der Verwendung: Kafka darf Aufgaben puffern, aber ein gelöstes Token sollte nicht minutenlang im captcha-results-Topic warten, bevor es abgesendet wird. Bei großem Rückstand ist es besser, die Aufgabe neu zu lösen als ein abgelaufenes Token zu übermitteln.

Kann ich verschiedene CAPTCHA-Typen über dasselbe Topic verarbeiten?

Ja. Jede Aufgabe trägt ihr eigenes method-Feld, sodass reCAPTCHA v2/v3, Cloudflare Turnstile, GeeTest v3 sowie Bild- und Raster-CAPTCHAs über dieselbe Pipeline laufen. Der Worker gibt den Wert unverändert an CaptchaAI weiter. hCaptcha und FunCaptcha unterstützt CaptchaAI nicht – solche Aufgaben gehören nicht in dieses Topic.

Wie stelle ich sicher, dass bei einem Worker-Absturz keine Aufgabe verloren geht?

Committen Sie den Offset erst nach der Verarbeitung (enable_auto_commit=False plus manueller commit()). So arbeitet die Pipeline nach dem At-least-once-Prinzip: Stürzt ein Worker vor dem Commit ab, wird die Aufgabe von einem anderen Worker erneut gelesen. Ergänzen Sie eine Idempotenzprüfung auf task_id, um daraus entstehende Doppelverarbeitungen abzufangen.

Verwandte Leitfäden

  • RabbitMQ als Message-Queue mit CaptchaAI verbinden
  • Verteilte Verarbeitung mit einer Redis-Warteschlange
  • Aufgaben mit NATS Messaging verteilen
Kommentare sind für diesen Artikel deaktiviert.