Streaming heißt hier: Ihre Pipeline verarbeitet jedes gelöste CAPTCHA in dem Moment, in dem die API es zurückgibt – und nicht erst, wenn der letzte Task des Batches fertig ist. Zwei Kennzahlen ändern sich dadurch sofort. Die Zeit bis zum ersten verwertbaren Token sinkt von der Laufzeit des gesamten Batches auf die der schnellsten Einzelaufgabe, und der Speicherbedarf bleibt konstant, statt mit der Batch-Größe zu wachsen.
Dieser Leitfaden zeigt das Muster zweimal: als asynchronen Generator in Python und als EventEmitter in Node.js – und klärt dabei, wie viele Aufgaben Sie überhaupt gleichzeitig laufen lassen dürfen.
Warum das erste Ergebnis mehr zählt als das letzte
Ein CAPTCHA-Batch wird nie gleichmäßig fertig: Manche Abfragen sind nach Sekunden beantwortet, andere brauchen ein Vielfaches, ein kleiner Rest läuft ins Timeout. Solange Ihr Code auf asyncio.gather() oder Promise.all() wartet, liegen fertige Token ungenutzt herum.
| Ansatz | Zeit bis zum ersten Ergebnis | Speicher | Pipeline-Latenz |
|---|---|---|---|
Auf alle warten (gather / Promise.all) |
nach der langsamsten Aufgabe | alle Ergebnisse im RAM | hoch |
| Streamen, sobald gelöst | nach der schnellsten Aufgabe | ein Ergebnis nach dem anderen | niedrig |
| Micro-Batch (Blöcke à 10) | nach dem ersten Block | 10 Ergebnisse gleichzeitig | mittel |
Der Micro-Batch ist der Mittelweg, wenn die nachgelagerte Verarbeitung ohnehin blockweise schreibt – etwa bei Datenbank-Inserts.
Wie viele Aufgaben Sie gleichzeitig streamen dürfen
Am Parallelitätslimit ändert Streaming nichts. CaptchaAI rechnet Thread-basiert ab: Ein Thread ist eine laufende CAPTCHA-Abfrage, und jeder Tarif enthält unbegrenzt viele Lösungen pro Thread im Abrechnungsmonat. Entscheidend ist damit nicht, wie viele CAPTCHAs Sie im Monat lösen, sondern wie viele gleichzeitig unterwegs sind.
Setzen Sie max_concurrent deshalb nie höher als das Thread-Kontingent Ihres Tarifs. STANDARD (30 $/Monat, 15 Threads) deckt genau die 15 parallelen Aufgaben aus den Beispielen unten ab, ADVANCE (90 $/Monat, 50 Threads) entsprechend mehr. Ein höheres Limit im Code bringt keinen zusätzlichen Durchsatz. Preise verstehen sich in US-Dollar.
Python: asynchroner Generator, der jedes Ergebnis sofort liefert
Das Muster besteht aus drei Bausteinen. submit_task reicht eine Aufgabe bei in.php ein, poll_task fragt res.php im Fünf-Sekunden-Takt ab, bis ein Token oder ein Fehlercode zurückkommt, und stream_results gibt jedes fertige Ergebnis per yield weiter. Die Ergebnisse erscheinen in der Reihenfolge ihrer Fertigstellung – deshalb trägt jedes seinen ursprünglichen index mit sich.
import asyncio
import aiohttp
import time
API_KEY = "YOUR_API_KEY"
SUBMIT_URL = "https://ocr.captchaai.com/in.php"
RESULT_URL = "https://ocr.captchaai.com/res.php"
async def submit_task(session, task_data):
"""Submit a single CAPTCHA task."""
params = {
"key": API_KEY,
"method": task_data.get("method", "userrecaptcha"),
"json": 1,
}
if params["method"] == "userrecaptcha":
params["googlekey"] = task_data["sitekey"]
params["pageurl"] = task_data["pageurl"]
elif params["method"] == "turnstile":
params["sitekey"] = task_data["sitekey"]
params["pageurl"] = task_data["pageurl"]
async with session.post(SUBMIT_URL, data=params) as resp:
result = await resp.json(content_type=None)
if result.get("status") != 1:
return None, result.get("request", "unknown")
return result["request"], None
async def poll_task(session, task_id, timeout=300):
"""Poll until solved or timeout."""
start = time.monotonic()
while time.monotonic() - start < timeout:
await asyncio.sleep(5)
params = {"key": API_KEY, "action": "get", "id": task_id, "json": 1}
async with session.get(RESULT_URL, params=params) as resp:
result = await resp.json(content_type=None)
if result.get("request") == "CAPCHA_NOT_READY":
continue
if result.get("status") == 1:
return result["request"], None
return None, result.get("request", "unknown")
return None, "TIMEOUT"
async def solve_one(session, index, task_data, semaphore):
"""Solve a single task within concurrency limits."""
async with semaphore:
start = time.monotonic()
task_id, error = await submit_task(session, task_data)
if error:
return {"index": index, "status": "failed", "error": error, "time": 0}
token, error = await poll_task(session, task_id)
elapsed = time.monotonic() - start
if token:
return {"index": index, "status": "solved", "token": token, "time": round(elapsed, 1)}
return {"index": index, "status": "failed", "error": error, "time": round(elapsed, 1)}
async def stream_results(tasks, max_concurrent=20):
"""
Async generator that yields each result as it completes.
Results arrive in completion order, not submission order.
"""
semaphore = asyncio.Semaphore(max_concurrent)
async with aiohttp.ClientSession() as session:
pending = set()
for i, task in enumerate(tasks):
coro = solve_one(session, i, task, semaphore)
pending.add(asyncio.ensure_future(coro))
while pending:
done, pending = await asyncio.wait(pending, return_when=asyncio.FIRST_COMPLETED)
for future in done:
yield future.result()
async def main():
tasks = [
{"sitekey": "SITE_KEY", "pageurl": f"https://example.com/page{i}"}
for i in range(50)
]
solved = 0
failed = 0
async for result in stream_results(tasks, max_concurrent=15):
# Process each result immediately
if result["status"] == "solved":
solved += 1
print(f" [{solved + failed}/{len(tasks)}] Task {result['index']} SOLVED in {result['time']}s")
# Use token immediately — don't wait for batch
# await submit_form(result["token"])
# await save_to_database(result)
else:
failed += 1
print(f" [{solved + failed}/{len(tasks)}] Task {result['index']} FAILED: {result['error']}")
print(f"\nDone: {solved} solved, {failed} failed")
asyncio.run(main())
Die einzige zusätzliche Abhängigkeit ist aiohttp:
pip install aiohttp
Drei Details entscheiden über die Stabilität unter Last:
asyncio.wait(..., return_when=FIRST_COMPLETED)ist der eigentliche Streaming-Kern: Die Schleife wacht auf, sobald ein Future fertig ist, und legt sich mit dem Rest wieder schlafen.- Die Semaphore begrenzt die offenen Abfragen. Alle Aufgaben entstehen sofort als Coroutine, unterwegs sind aber nur
max_concurrentdavon. poll_tasktrennt Timeout und echte Fehlercodes. Wer beides vermischt, sieht später nicht mehr, ob eine Aufgabe zu langsam war oder falsch eingereicht wurde.
Node.js: EventEmitter statt Warten auf Promise.all
In Node.js übernimmt ein EventEmitter die Rolle des asynchronen Generators. Jede fertige Aufgabe löst ein result-Event aus, der Listener verarbeitet sie sofort, done markiert das Ende des Laufs. processNext() hält konstant so viele Anfragen offen, wie maxConcurrent erlaubt – eine neue rückt erst nach, wenn eine andere abgeschlossen ist.
const { EventEmitter } = require("events");
const API_KEY = "YOUR_API_KEY";
const SUBMIT_URL = "https://ocr.captchaai.com/in.php";
const RESULT_URL = "https://ocr.captchaai.com/res.php";
class CaptchaStream extends EventEmitter {
constructor(maxConcurrent = 15) {
super();
this.maxConcurrent = maxConcurrent;
this.active = 0;
this.queue = [];
this.total = 0;
this.completed = 0;
}
async submitAndPoll(index, taskData) {
const params = new URLSearchParams({
key: API_KEY,
method: taskData.method || "userrecaptcha",
googlekey: taskData.sitekey,
pageurl: taskData.pageurl,
json: "1",
});
const start = Date.now();
const submitResp = await (await fetch(SUBMIT_URL, { method: "POST", body: params })).json();
if (submitResp.status !== 1) {
return { index, status: "failed", error: submitResp.request, time: 0 };
}
const taskId = submitResp.request;
for (let i = 0; i < 60; i++) {
await new Promise((r) => setTimeout(r, 5000));
const url = `${RESULT_URL}?key=${API_KEY}&action=get&id=${taskId}&json=1`;
const poll = await (await fetch(url)).json();
if (poll.request === "CAPCHA_NOT_READY") continue;
const elapsed = ((Date.now() - start) / 1000).toFixed(1);
if (poll.status === 1) return { index, status: "solved", token: poll.request, time: elapsed };
return { index, status: "failed", error: poll.request, time: elapsed };
}
return { index, status: "failed", error: "TIMEOUT", time: ((Date.now() - start) / 1000).toFixed(1) };
}
async processNext() {
if (this.queue.length === 0 || this.active >= this.maxConcurrent) return;
const { index, taskData } = this.queue.shift();
this.active++;
try {
const result = await this.submitAndPoll(index, taskData);
this.emit("result", result);
} catch (err) {
this.emit("result", { index, status: "failed", error: err.message });
} finally {
this.active--;
this.completed++;
if (this.completed === this.total) {
this.emit("done");
} else {
this.processNext();
}
}
}
start(tasks) {
this.total = tasks.length;
this.queue = tasks.map((taskData, index) => ({ index, taskData }));
// Launch initial batch
const initial = Math.min(this.maxConcurrent, tasks.length);
for (let i = 0; i < initial; i++) {
this.processNext();
}
return this;
}
}
// Usage
const tasks = Array.from({ length: 50 }, (_, i) => ({
sitekey: "SITE_KEY",
pageurl: `https://example.com/page${i}`,
}));
const stream = new CaptchaStream(15);
let solved = 0, failed = 0;
stream.on("result", (result) => {
if (result.status === "solved") {
solved++;
console.log(`[${solved + failed}/${tasks.length}] Task ${result.index} SOLVED (${result.time}s)`);
// Use token immediately
// submitForm(result.token);
} else {
failed++;
console.log(`[${solved + failed}/${tasks.length}] Task ${result.index} FAILED: ${result.error}`);
}
});
stream.on("done", () => {
console.log(`\nComplete: ${solved} solved, ${failed} failed`);
});
stream.start(tasks);
Praxisbeispiel: nächtlicher Formular-Check aus der GitLab-CI
Ein typischer Aufbau aus dem DACH-Raum: Ein Team betreibt die Staging-Instanz seines Shopware-Shops auf einem Hetzner-VPS und prüft nachts per GitLab-CI-Job rund 1.200 Formularstrecken unter https://staging.example-app.test. Jede Strecke enthält ein reCAPTCHA v2 und braucht ein frisches Token.
Sammelt der Job erst alles ein, steht die Pipeline bis zur letzten Abfrage – und der Runner hält die komplette Ergebnisliste im Speicher. Mit Streaming startet der erste Formular-Submit, sobald das erste Token vorliegt: Fehler werden während des Laufs sichtbar, und ein Abbruch nach 40 Minuten hinterlässt trotzdem mehrere hundert geprüfte Strecken statt gar keiner.
Noch ein Hinweis zur Protokollierung: Enthalten Ihre Task-Daten vollständige Ziel-URLs oder IP-Adressen aus Proxy-Logs, berührt das die DSGVO. Ein Handler, der nur Index, Status und Lösungszeit schreibt, ist datensparsamer – und reicht zur Auswertung meist aus.
Streamen oder doch sammeln?
| Szenario | Empfehlung |
|---|---|
| Formulare mit frisch gelösten Token absenden | streamen – jedes Formular geht raus, sobald sein Token da ist |
| CSV- oder Parquet-Export des Laufs | sammeln – einmal schreiben, wenn der Batch durch ist |
| Dashboard mit Live-Fortschritt | streamen – Oberfläche bei jedem result-Event aktualisieren |
| Aufgaben mit Abhängigkeiten untereinander | sammeln – danach in definierter Reihenfolge verarbeiten |
| Große Läufe ab 1.000 Aufgaben | streamen – hält den Spitzenspeicher konstant |
Das stärkste Argument für Streaming ist die Token-Lebensdauer: reCAPTCHA-Token sind nur rund 120 Sekunden gültig. Wer erst nach einem zwanzigminütigen Lauf mit der Verarbeitung beginnt, hält für den Anfang des Batches längst abgelaufene Token in der Hand.
Reihenfolge, Speicher und Backpressure
An drei Stellen kippen Streaming-Pipelines in der Praxis.
Reihenfolge. Ergebnisse treffen nach Fertigstellung ein. Braucht die nachgelagerte Verarbeitung die ursprüngliche Sortierung, puffern Sie sie in einem Dictionary nach index und geben Sie zusammenhängende Blöcke frei, sobald die Lücke geschlossen ist – dasselbe Prinzip wie beim Zusammensetzen von TCP-Paketen.
Speicher. Streaming spart nur RAM, wenn der Handler das Ergebnis auch verwirft – ein results.append(result) im Listener macht den Vorteil zunichte.
Backpressure. Schreibt Ihr Handler jedes Token synchron in eine Datenbank, wird er selbst zum Flaschenhals. Wie Sie dann drosseln, statt die Warteschlange volllaufen zu lassen, zeigt der Leitfaden zu Backpressure in CAPTCHA-Warteschlangen.
Fehlerbehebung
| Symptom | Ursache | Lösung |
|---|---|---|
| Ergebnisse kommen in zufälliger Reihenfolge | erwartetes Verhalten – Schnellste zuerst | über result.index zurückmappen |
| Das erste Ergebnis dauert sehr lange | alle Aufgaben gleichzeitig eingereicht | Einreichung über Semaphore bzw. maxConcurrent staffeln |
Node.js meldet MaxListenersExceeded |
zu viele Listener am selben Stream | pro Event-Typ genau einen Listener registrieren |
| Der asynchrone Generator hängt | ein Future im pending-Set wird nie fertig |
Timeout in poll_task erzwingen |
| Token wird geliefert, aber vom Formular abgelehnt | sitekey, pageurl oder Session-Kontext passen nicht |
Parameter neu erfassen, Token in derselben Sitzung einsetzen |
Häufige Fragen
Wie viele CAPTCHAs kann ich gleichzeitig streamen?
So viele, wie Ihr Tarif Threads hat. Ein Thread entspricht einer laufenden Abfrage; ist sie beantwortet, nimmt er sofort die nächste. Die Beispiele oben nutzen 15 gleichzeitige Aufgaben und passen auf STANDARD (30 $/Monat, 15 Threads).
Was passiert, wenn ein Token abläuft, bevor ich es einsetze?
Dann weist das Zielformular es zurück. reCAPTCHA-Token sind nur rund 120 Sekunden gültig – deshalb lösen Sie jedes Token direkt vor der Übermittlung und verwenden es sofort, statt Vorräte anzulegen.
Lässt sich ein abgebrochener Lauf fortsetzen?
Ja. Schreiben Sie jedes Ergebnis beim Eintreffen als JSON-Zeile in eine Checkpoint-Datei. Beim Neustart laden Sie die Datei, filtern die erledigten Indizes heraus und reichen nur den Rest erneut ein.
Funktioniert das Muster mit gemischten CAPTCHA-Typen?
Ja. Jede Aufgabe trägt ihre eigene method – im Python-Beispiel userrecaptcha für reCAPTCHA v2 und turnstile für Cloudflare Turnstile. Die Streaming-Schleife gibt weiter, was fertig wird; GeeTest v3 und Bild-CAPTCHAs hängen Sie nach demselben Schema ein.