Um lote de 500 tarefas CAPTCHA é concluído de maneira irregular: algumas são resolvidas em 8 segundos, outras levam 45. Esperar que cada tarefa termine antes de processar os resultados desperdiça o tempo entre a primeira e a última solução. O streaming permite que seu pipeline downstream consuma cada resultado no momento em que ele chega.
Streaming versus processamento em lote
| Abordagem | Hora do primeiro resultado | Memória | Latência do pipeline |
|---|---|---|---|
| Espere por todos | Depois da tarefa mais lenta | Todos os resultados na memória | Alto |
| Transmitir como resolvido | Depois da tarefa com menor latência | Um resultado de cada vez | Baixo |
| Microlote (pedaços de 10) | Depois do primeiro pedaço | 10 resultados por vez | Médio |
Python: gerador assíncrono para resultados de streaming
Usando asyncio e aiohttp, cada solução produz resultados imediatos por meio de um gerador assíncrono:
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())
Instale dependências:
pip install aiohttp
JavaScript: padrão de streaming EventEmitter
O Node.js usa uma abordagem orientada a eventos – emita cada resultado conforme ele resolve:
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);
Quando usar streaming versus coletar tudo
| Cenário | Abordagem |
|---|---|
| Envios de formulários usando tokens | Stream – envie cada formulário assim que o token chegar |
| Exportação CSV de todos os resultados | Coletar tudo – escreva uma vez quando o lote for concluído |
| Painel com progresso ao vivo | Stream – atualiza a UI em cada evento de resultado |
| Lote com dependências entre tarefas | Colete tudo - processe em ordem após a conclusão |
| Grandes lotes (mais de 1.000) | Stream – reduz o pico de uso de memória |
Solução de problemas
| Problema | Causa | Correção |
|---|---|---|
| Os resultados chegam em ordem aleatória | Normal – o streaming produz o rendimento com menor latência primeiro | Use result.index para mapear de volta à tarefa original |
| A memória ainda cresce durante a transmissão | Armazenando todos os resultados em array | Processar e descartar resultados no manipulador |
| O primeiro resultado demora muito | Todas as tarefas enviadas simultaneamente | Escalonar envios com semáforo ou limite de simultaneidade |
| Aviso do EventEmitter: MaxListenersExceeded | Muitos ouvintes na transmissão | Use setMaxListeners() ou garanta um ouvinte por tipo de evento |
| Gerador assíncrono trava | Tarefa não resolvida em conjunto pendente | Adicione tempo limite a poll_task; garantir que todos os futuros estejam completos ou com erro |
Perguntas frequentes
O streaming aumenta as chamadas de API em comparação com o lote?
Não - o mesmo número de chamadas de envio e votação acontece de qualquer maneira. O streaming só muda quando seu aplicativo processa cada resultado, e não quantas chamadas de API são feitas.
Como mantenho a ordem das tarefas durante o streaming?
Cada resultado carrega seu index original. Se a ordem for importante para o processamento downstream, o buffer resulta em uma estrutura classificada e libera execuções contíguas (como a remontagem de pacotes TCP).
Posso combinar streaming com checkpoint?
Sim. Anexe cada resultado a um arquivo de ponto de verificação assim que ele chegar. Ao retomar, carregue o ponto de verificação, filtre os índices concluídos e processe novamente apenas as tarefas restantes.
Artigos relacionados
- Processamento de Captcha de streaming de Kafka Captchaai
- Processamento prioritário de Captcha em lote de fila
Próximas etapas
Processe soluções CAPTCHA no momento em que elas chegam -obtenha sua chave API CaptchaAIe construir pipelines de streaming.
Guias relacionados: