Tutorials

Erstellen einer CAPTCHA-Lösungswarteschlange in Python mit CaptchaAI

Sobald Sie mehr als ein CAPTCHA pro Sekunde verarbeiten, ist eine Warteschlange die richtige Antwort: Sie übermitteln alle Abfragen sofort, fragen die Ergebnisse parallel ab und binden die Zahl der gleichzeitigen Worker an Ihr Thread-Kontingent bei CaptchaAI. Dieser Artikel zeigt vier Muster in Python – von einer schlanken Threading-Warteschlange bis zur Prioritätswarteschlange – und wie Sie die Parallelität sauber dimensionieren, statt die API zu überlasten.


Wann sich eine Warteschlange lohnt

Ein CAPTCHA sequenziell zu lösen bedeutet: absenden, dann fünf bis fünfzehn Sekunden auf das Ergebnis warten, erst danach das nächste starten. Bei 100 Seiten summiert sich diese Wartezeit zu Minuten, in denen Ihr Prozess nichts tut. Eine Warteschlange dreht das um – sie hält viele Lösungen gleichzeitig „in der Luft".

Ein Warteschlangensystem übernimmt dabei fünf Aufgaben:

  • Es übermittelt alle CAPTCHAs sofort, ohne auf das vorherige Ergebnis zu warten.
  • Es fragt mehrere Aufgaben-IDs parallel ab (Polling).
  • Es wiederholt fehlgeschlagene Abfragen automatisch.
  • Es begrenzt die Parallelität, damit Sie das Anfragelimit der API einhalten.
  • Es liefert Fortschrittsanzeige und Callbacks für weiterverarbeitende Logik.

Der entscheidende Hebel ist die Zahl der gleichzeitigen Worker. Zu wenige verschenken Durchsatz, zu viele erzeugen ERROR_NO_SLOT_AVAILABLE. Wie Sie diesen Wert an Ihren Plan koppeln, klärt der Abschnitt weiter unten.


Threading-Warteschlange: der einfache Einstieg

Für bestehenden synchronen Code – etwa ein klassisches Scraping-Skript mit requests – ist eine Threading-Warteschlange der pragmatischste Weg. Worker-Threads ziehen Aufgaben aus einer Queue, lösen sie über CaptchaAI und legen das Ergebnis in einer zweiten Warteschlange ab. Die submit-Methode nimmt reCAPTCHA-, Turnstile- oder Bild-CAPTCHA-Aufgaben entgegen; die Methodenkennung (userrecaptcha, turnstile und so weiter) entscheidet über den CAPTCHA-Typ.

import time
import threading
import requests
from queue import Queue, Empty

API_KEY = "YOUR_API_KEY"


class CaptchaQueue:
    """Thread-based CAPTCHA solving queue."""

    def __init__(self, api_key, max_workers=10):
        self.api_key = api_key
        self.task_queue = Queue()
        self.result_queue = Queue()
        self.max_workers = max_workers
        self.workers = []

    def submit(self, method, callback=None, **params):
        """Add a CAPTCHA task to the queue."""
        task = {
            "method": method,
            "params": params,
            "callback": callback,
        }
        self.task_queue.put(task)

    def start(self):
        """Start worker threads."""
        for _ in range(self.max_workers):
            t = threading.Thread(target=self._worker, daemon=True)
            t.start()
            self.workers.append(t)

    def wait(self):
        """Wait for all tasks to complete."""
        self.task_queue.join()

    def get_results(self):
        """Get all available results."""
        results = []
        while not self.result_queue.empty():
            try:
                results.append(self.result_queue.get_nowait())
            except Empty:
                break
        return results

    def _worker(self):
        while True:
            try:
                task = self.task_queue.get(timeout=1)
            except Empty:
                continue

            try:
                result = self._solve(task["method"], **task["params"])
                entry = {"status": "solved", "result": result, "task": task}
                self.result_queue.put(entry)
                if task["callback"]:
                    task["callback"](result)
            except Exception as e:
                entry = {"status": "error", "error": str(e), "task": task}
                self.result_queue.put(entry)
            finally:
                self.task_queue.task_done()

    def _solve(self, method, **params):
        submit = requests.post("https://ocr.captchaai.com/in.php", data={
            "key": self.api_key, "method": method, "json": 1, **params,
        }, timeout=30).json()

        if submit.get("status") != 1:
            raise Exception(f"Submit error: {submit.get('request')}")

        task_id = submit["request"]
        for _ in range(30):
            time.sleep(5)
            result = requests.get("https://ocr.captchaai.com/res.php", params={
                "key": self.api_key, "action": "get", "id": task_id, "json": 1,
            }, timeout=30).json()
            if result.get("status") == 1:
                return result["request"]
            if result.get("request") == "ERROR_CAPTCHA_UNSOLVABLE":
                raise Exception("CAPTCHA unsolvable")
        raise TimeoutError("Solve timed out")


# Usage
queue = CaptchaQueue(API_KEY, max_workers=5)
queue.start()

# Submit multiple CAPTCHAs
urls_and_sitekeys = [
    ("https://example.com/page1", "SITEKEY_1"),
    ("https://example.com/page2", "SITEKEY_2"),
    ("https://example.com/page3", "SITEKEY_3"),
]

for url, sitekey in urls_and_sitekeys:
    queue.submit("userrecaptcha", googlekey=sitekey, pageurl=url)

queue.wait()
results = queue.get_results()
print(f"Solved {len(results)} CAPTCHAs")
for r in results:
    print(f"  {r['status']}: {r.get('result', r.get('error', ''))[:50]}")

Jeder Worker kapselt Übermittlung und Polling; task_done() im finally-Block sorgt dafür, dass wait() zuverlässig zurückkehrt, auch wenn eine Lösung fehlschlägt. Die Methodenkennung userrecaptcha deckt dabei alle reCAPTCHA-v2-Varianten ab.


Asyncio-Warteschlange für I/O-lastige Workloads

CAPTCHA-Lösen ist fast reine Wartezeit auf das Netzwerk – der klassische I/O-bound-Fall, in dem Asyncio glänzt. Statt für jeden Worker einen Betriebssystem-Thread zu binden, verwalten Sie Hunderte gleichzeitiger Coroutinen in einer einzigen Event-Loop. Ein asyncio.Semaphore begrenzt dabei, wie viele Lösungen wirklich gleichzeitig laufen.

import asyncio
import aiohttp

API_KEY = "YOUR_API_KEY"


class AsyncCaptchaQueue:
    """Async CAPTCHA solving queue with concurrency control."""

    def __init__(self, api_key, max_concurrent=10):
        self.api_key = api_key
        self.semaphore = asyncio.Semaphore(max_concurrent)
        self.results = []

    async def solve_batch(self, tasks):
        """Solve a batch of CAPTCHA tasks concurrently."""
        coros = [self._solve_task(task) for task in tasks]
        self.results = await asyncio.gather(*coros, return_exceptions=True)
        return self.results

    async def _solve_task(self, task):
        async with self.semaphore:
            return await self._solve(task["method"], **task["params"])

    async def _solve(self, method, **params):
        async with aiohttp.ClientSession() as session:
            # Submit
            async with session.post("https://ocr.captchaai.com/in.php", data={
                "key": self.api_key, "method": method, "json": 1, **params,
            }) as resp:
                data = await resp.json(content_type=None)
                if data.get("status") != 1:
                    raise Exception(f"Submit error: {data.get('request')}")
                task_id = data["request"]

            # Poll
            for _ in range(30):
                await asyncio.sleep(5)
                async with session.get("https://ocr.captchaai.com/res.php", params={
                    "key": self.api_key, "action": "get", "id": task_id, "json": 1,
                }) as resp:
                    result = await resp.json(content_type=None)
                    if result.get("status") == 1:
                        return result["request"]
                    if result.get("request") == "ERROR_CAPTCHA_UNSOLVABLE":
                        raise Exception("CAPTCHA unsolvable")

            raise TimeoutError("Solve timed out")


# Usage
async def main():
    queue = AsyncCaptchaQueue(API_KEY, max_concurrent=5)

    tasks = [
        {"method": "userrecaptcha", "params": {"googlekey": f"SITEKEY_{i}", "pageurl": f"https://example.com/page{i}"}}
        for i in range(10)
    ]

    results = await queue.solve_batch(tasks)
    for i, result in enumerate(results):
        if isinstance(result, Exception):
            print(f"Task {i}: ERROR — {result}")
        else:
            print(f"Task {i}: {result[:50]}...")


asyncio.run(main())

Mit return_exceptions=True bricht asyncio.gather nicht beim ersten Fehler ab – jede Aufgabe liefert entweder ein Token oder eine Exception, die Sie einzeln auswerten. Dieses Muster ist die richtige Wahl für neue Projekte und skaliert deutlich weiter als eine Thread-pro-Worker-Architektur.


Producer-Consumer-Muster für kontinuierliches Scraping

Bei einem laufenden Crawler kennen Sie die Liste der Seiten nicht vorab – neue URLs tauchen erst auf, während Sie andere abarbeiten. Hier passt das Producer-Consumer-Muster: Ein Producer schiebt neu entdeckte Aufgaben in eine asyncio.Queue, mehrere Consumer lösen sie parallel. Ein None-Sentinel signalisiert das Ende sauber pro Consumer.

import asyncio
import aiohttp

API_KEY = "YOUR_API_KEY"


class ProducerConsumerQueue:
    """Continuous CAPTCHA solving with producer-consumer pattern."""

    def __init__(self, api_key, queue_size=100, num_consumers=5):
        self.api_key = api_key
        self.queue = asyncio.Queue(maxsize=queue_size)
        self.num_consumers = num_consumers
        self.solved_count = 0
        self.error_count = 0
        self.running = True

    async def produce(self, tasks):
        """Producer: feed CAPTCHA tasks into the queue."""
        for task in tasks:
            await self.queue.put(task)
        # Signal consumers to stop
        for _ in range(self.num_consumers):
            await self.queue.put(None)

    async def consume(self, result_handler):
        """Consumer: solve CAPTCHAs and call result handler."""
        async with aiohttp.ClientSession() as session:
            while True:
                task = await self.queue.get()
                if task is None:
                    self.queue.task_done()
                    break

                try:
                    result = await self._solve(session, task["method"], **task["params"])
                    self.solved_count += 1
                    if result_handler:
                        await result_handler(task, result)
                except Exception as e:
                    self.error_count += 1
                    print(f"Error: {e}")
                finally:
                    self.queue.task_done()

    async def run(self, tasks, result_handler=None):
        """Run the producer-consumer pipeline."""
        # Start producer
        producer = asyncio.create_task(self.produce(tasks))

        # Start consumers
        consumers = [
            asyncio.create_task(self.consume(result_handler))
            for _ in range(self.num_consumers)
        ]

        # Wait for everything to finish
        await producer
        await asyncio.gather(*consumers)

        print(f"Complete: {self.solved_count} solved, {self.error_count} errors")

    async def _solve(self, session, method, **params):
        async with session.post("https://ocr.captchaai.com/in.php", data={
            "key": self.api_key, "method": method, "json": 1, **params,
        }) as resp:
            data = await resp.json(content_type=None)
            if data.get("status") != 1:
                raise Exception(f"Submit: {data.get('request')}")
            task_id = data["request"]

        for _ in range(30):
            await asyncio.sleep(5)
            async with session.get("https://ocr.captchaai.com/res.php", params={
                "key": self.api_key, "action": "get", "id": task_id, "json": 1,
            }) as resp:
                result = await resp.json(content_type=None)
                if result.get("status") == 1:
                    return result["request"]
        raise TimeoutError("Timed out")


# Usage
async def handle_result(task, token):
    url = task["params"]["pageurl"]
    print(f"Solved for {url}: {token[:30]}...")


async def main():
    queue = ProducerConsumerQueue(API_KEY, num_consumers=5)

    tasks = [
        {"method": "userrecaptcha", "params": {"googlekey": f"SITEKEY_{i}", "pageurl": f"https://example.com/page{i}"}}
        for i in range(20)
    ]

    await queue.run(tasks, result_handler=handle_result)


asyncio.run(main())

Die Warteschlangengröße (queue_size) wirkt als Puffer: Produziert der Crawler schneller, als die Consumer lösen können, blockiert put und bremst den Producer automatisch aus – ein eingebauter Backpressure-Mechanismus, der Ihren Arbeitsspeicher schont.


Prioritätswarteschlange: zeitkritische Abfragen zuerst

Nicht jede Aufgabe ist gleich dringend. Eine asyncio.PriorityQueue arbeitet Aufgaben nach einem Prioritätswert ab (kleinere Zahl = höhere Priorität), sodass zeitkritische Abfragen nicht hinter einem Stapel unwichtiger Seiten warten müssen.

import asyncio
from dataclasses import dataclass, field

API_KEY = "YOUR_API_KEY"


@dataclass(order=True)
class PriorityTask:
    priority: int
    task: dict = field(compare=False)


class PriorityCaptchaQueue:
    """CAPTCHA queue with priority levels."""

    def __init__(self, api_key, num_workers=5):
        self.api_key = api_key
        self.queue = asyncio.PriorityQueue()
        self.num_workers = num_workers
        self.results = {}

    async def submit(self, task_id, method, priority=5, **params):
        """Submit with priority (lower number = higher priority)."""
        await self.queue.put(PriorityTask(
            priority=priority,
            task={"id": task_id, "method": method, "params": params},
        ))

    async def process(self):
        """Process all queued tasks by priority."""
        workers = [asyncio.create_task(self._worker()) for _ in range(self.num_workers)]

        # Wait for queue to drain
        await self.queue.join()

        # Cancel workers
        for w in workers:
            w.cancel()

        return self.results

    async def _worker(self):
        import aiohttp
        async with aiohttp.ClientSession() as session:
            while True:
                item = await self.queue.get()
                task = item.task
                try:
                    result = await self._solve(session, task["method"], **task["params"])
                    self.results[task["id"]] = {"status": "solved", "token": result}
                except Exception as e:
                    self.results[task["id"]] = {"status": "error", "error": str(e)}
                finally:
                    self.queue.task_done()

    async def _solve(self, session, method, **params):
        import aiohttp
        async with session.post("https://ocr.captchaai.com/in.php", data={
            "key": self.api_key, "method": method, "json": 1, **params,
        }) as resp:
            data = await resp.json(content_type=None)
            if data.get("status") != 1:
                raise Exception(data.get("request"))
            task_id = data["request"]

        for _ in range(30):
            await asyncio.sleep(5)
            async with session.get("https://ocr.captchaai.com/res.php", params={
                "key": self.api_key, "action": "get", "id": task_id, "json": 1,
            }) as resp:
                result = await resp.json(content_type=None)
                if result.get("status") == 1:
                    return result["request"]
        raise TimeoutError()


# Usage
async def main():
    pq = PriorityCaptchaQueue(API_KEY, num_workers=3)

    # High priority — checkout pages
    await pq.submit("checkout_1", "turnstile", priority=1, sitekey="KEY", pageurl="https://shop.com/checkout")

    # Normal priority — product pages
    for i in range(5):
        await pq.submit(f"product_{i}", "userrecaptcha", priority=5, googlekey="KEY", pageurl=f"https://shop.com/p/{i}")

    # Low priority — info pages
    for i in range(3):
        await pq.submit(f"info_{i}", "userrecaptcha", priority=10, googlekey="KEY", pageurl=f"https://shop.com/info/{i}")

    results = await pq.process()
    for task_id, result in results.items():
        print(f"{task_id}: {result['status']}")


asyncio.run(main())

Das Beispiel mischt bewusst Typen: ein Turnstile-CAPTCHA mit hoher Priorität, mehrere reCAPTCHA-Aufgaben mit normaler und niedriger Priorität. So verarbeitet dieselbe Warteschlange eine zeitkritische Formularseite vorrangig, während unwichtige Übersichtsseiten warten. Für Turnstile-Aufgaben genügt die Methodenkennung turnstile zusammen mit sitekey und pageurl.


Threads und Parallelität an Ihren Plan koppeln

Der häufigste Fehler beim Skalieren: mehr Worker starten, als der Plan hergibt. CaptchaAI rechnet pro gleichzeitigem Thread ab, nicht pro Lösung – jeder Plan enthält unbegrenzte Lösungen pro Thread im Abrechnungsmonat. Ein Thread ist ein CAPTCHA „in Bearbeitung"; sobald eine Lösung fertig ist, nimmt der Thread die nächste Aufgabe auf.

Praktisch heißt das: Ihre max_workers beziehungsweise max_concurrent sollten Ihr Thread-Kontingent nicht überschreiten. Zur Orientierung an den kanonischen Preisen (in US-Dollar):

  • BASIC (15 $/Monat, 5 Threads) – erste Tests und kleine Skripte
  • STANDARD (30 $/Monat, 15 Threads) – regelmäßiges Scraping mittlerer Größe
  • ADVANCE (90 $/Monat, 50 Threads) – parallele Pipelines mit dauerhaftem Durchsatz
  • PREMIUM (170 $/Monat, 100 Threads) – umfangreiche, kontinuierliche Workloads

Setzen Sie die Worker-Zahl an das Thread-Kontingent und lassen Sie etwas Reserve. Läuft der Crawler auf einem eigenen Server – im DACH-Raum etwa bei Hetzner oder netcup –, dann bündeln Sie mehrere Worker in einem Prozess, statt viele kleine Instanzen zu betreiben. Und da beim Web-Scraping IP-Adressen als personenbezogene Daten gelten, sollten Sie Ihre Datenflüsse und die Rechtsgrundlage nach DSGVO vorab prüfen.


Metriken und Monitoring

Ohne Kennzahlen tappen Sie im Dunkeln. Eine schlanke Metrik-Klasse erfasst Durchsatz, durchschnittliche Lösungszeit und Erfolgsquote – die Werte, an denen Sie ablesen, ob Ihre Parallelität passt.

import time
from dataclasses import dataclass, field


@dataclass
class QueueMetrics:
    submitted: int = 0
    solved: int = 0
    failed: int = 0
    total_solve_time: float = 0.0
    start_time: float = field(default_factory=time.time)

    @property
    def avg_solve_time(self):
        return self.total_solve_time / self.solved if self.solved else 0

    @property
    def success_rate(self):
        total = self.solved + self.failed
        return (self.solved / total * 100) if total else 0

    @property
    def throughput(self):
        elapsed = time.time() - self.start_time
        return self.solved / elapsed * 60 if elapsed > 0 else 0

    def report(self):
        return (
            f"Submitted: {self.submitted} | "
            f"Solved: {self.solved} | "
            f"Failed: {self.failed} | "
            f"Avg time: {self.avg_solve_time:.1f}s | "
            f"Success: {self.success_rate:.1f}% | "
            f"Throughput: {self.throughput:.0f}/min"
        )

Sinkt die Erfolgsquote während einer Lastspitze, ist das oft ein Zeichen für zu viele parallele Worker – ein Fall für den Wert aus dem vorigen Abschnitt.


Fehlerbehebung

Symptom Ursache Lösung
Warteschlange wächst, Aufgaben werden nicht fertig Zu viele Worker überlasten die API max_workers / max_concurrent reduzieren
ERROR_NO_SLOT_AVAILABLE Thread-Limit des Plans erreicht Verzögerung zwischen Übermittlungen einbauen oder Worker reduzieren
Aufgaben hängen dauerhaft in der Warteschlange Worker-Thread ist an einer Exception gestorben Worker-Schleife in try/except kapseln
Der Arbeitsspeicher wächst mit der Zeit Ergebnisse werden nicht abgeholt get_results() regelmäßig aufrufen
Asyncio-Warteschlange blockiert Ein await fehlt Sicherstellen, dass alle asynchronen Aufrufe abgewartet werden

Häufige Fragen

Wie viele Worker passen zu meinem Thread-Kontingent?

So viele, wie Ihr Thread-Kontingent zulässt – mit etwas Reserve. Bei ADVANCE (90 $/Monat, 50 Threads) sind rund 40 bis 50 parallele Worker realistisch. Steigt ERROR_NO_SLOT_AVAILABLE, ist das Limit erreicht.

Was bedeutet ERROR_NO_SLOT_AVAILABLE?

Sie haben mehr gleichzeitige Aufgaben übermittelt, als Ihr Plan an Threads bietet. Bauen Sie eine kurze Verzögerung zwischen den Übermittlungen ein oder senken Sie die Worker-Zahl, bis der Fehler verschwindet.

Wie oft sollte ich das Ergebnis abfragen?

Ein Polling-Intervall von etwa 5 Sekunden ist ein guter Startwert – so wie in den Beispielen. Kürzere Intervalle erhöhen nur die Anfragelast, ohne die Lösung zu beschleunigen.

Kann ich verschiedene CAPTCHA-Typen in einer Warteschlange mischen?

Ja. Die Methodenkennung pro Aufgabe (userrecaptcha, turnstile, geetest und weitere) bestimmt den Typ, sodass reCAPTCHA, Turnstile und Bild-CAPTCHAs problemlos in derselben Warteschlange laufen.

Threading oder Asyncio – was ist die bessere Wahl?

Für neue Projekte Asyncio, weil es die netzwerklastige Lösung effizienter parallelisiert. Threading lohnt sich, wenn Sie in bestehenden synchronen Code integrieren.


Fazit

Eine CAPTCHA-Warteschlange entkoppelt die Übermittlung von der Ergebnisabfrage und ermöglicht so paralleles Lösen mit CaptchaAI. Wählen Sie Threading für synchronen Bestandscode, Asyncio für moderne Projekte und das Producer-Consumer-Muster für kontinuierliche Crawler – und dimensionieren Sie die Worker-Zahl immer an Ihrem Thread-Kontingent. GeeTest-v3- oder Bild-CAPTCHAs binden Sie nach demselben Prinzip ein – nur die Methodenkennung ändert sich.

Verwandte Artikel

Kommentare sind für diesen Artikel deaktiviert.