500 CAPTCHA görevinden oluşan bir grup düzensiz bir şekilde tamamlanıyor; bazıları 8 saniyede, bazıları ise 45 saniyede çözülüyor. Sonuçları işlemeden önce her görevin bitmesini beklemek, ilk ve son çözüm arasındaki zamanı boşa harcar. Akış, aşağı akış işlem hattınızın her sonucu geldiği anda tüketmesine olanak tanır.
Akış ve Sonra Toplu İşlem Karşılaştırması
| Yaklaşım | İlk Sonuca Kadar Zaman | Bellek | Boru Hattı Gecikmesi |
|---|---|---|---|
| Hepsini bekle | En yavaş görevden sonra | Tüm sonuçlar hafızada | Yüksek |
| Çözülmüş olarak yayınla | En hızlı görevden sonra | Tek seferde tek sonuç | Düşük |
| Mikro parti (10'luk parçalar) | İlk parçadan sonra | Tek seferde 10 sonuç | Orta |
Python: Sonuç Akışı için Eşzamansız Oluşturucu
asyncio ve aiohttp kullanıldığında, her çözüm bir eşzamansız oluşturucu aracılığıyla anında sonuç verir:
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())
Bağımlılıkları yükleyin:
pip install aiohttp
JavaScript: EventEmitter Akış Kalıbı
Node.js, olaya dayalı bir yaklaşım kullanır; her sonucu çözüldükçe yayınlar:
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);
Akış ve Tümünü Toplama Ne Zaman Kullanılmalı?
| Senaryo | Yaklaşım |
|---|---|
| Belirteçleri kullanarak form gönderimleri | Akış – her formu jeton gelir gelmez gönderin |
| Tüm sonuçların CSV dışa aktarımı | Tümünü topla; grup tamamlandığında bir kez yazın |
| Canlı ilerlemeyi gösteren kontrol paneli | Akış – her sonuç etkinliğinde kullanıcı arayüzünü güncelleyin |
| Görevler arası bağımlılıklara sahip toplu iş | Tamamlandıktan sonra tümünü sırayla toplayın |
| Büyük partiler (1.000+) | Akış – en yüksek bellek kullanımını azaltın |
Sorun giderme
| Sorun | Sebep | Düzeltme |
|---|---|---|
| Sonuçlar rastgele sırada geliyor | Normal — yayın akışı en hızlı şekilde önce gelir | Orijinal göreve geri dönmek için result.index kullanın |
| Bellek akış sırasında büyümeye devam ediyor | Tüm sonuçların dizide saklanması | Sonuçları işleyicide işleyin ve atın |
| İlk sonuç çok uzun sürüyor | Tüm görevler aynı anda gönderildi | Semafor veya eşzamanlılık sınırıyla gönderimleri kademelendirin |
| EventEmitter uyarısı: MaxListenersExceeded | Yayında çok fazla dinleyici var | setMaxListeners() kullanın veya etkinlik türü başına bir dinleyici sağlayın |
| Zaman uyumsuz oluşturucu kilitleniyor | Bekleyen kümedeki çözülmemiş görev | poll_task'ye zaman aşımı ekleyin; tüm vadeli işlemlerin tamamlandığından veya hatalı olduğundan emin olun |
SSS
Akış, toplu işleme kıyasla API çağrılarını artırıyor mu?
Hayır; her iki durumda da aynı sayıda gönderme ve anket çağrısı gerçekleşir. Akış, kaç API çağrısının yapıldığını değil, yalnızca uygulamanız her sonucu işlediğinde değişir.
Akış sırasında görev sırasını nasıl koruyabilirim?
Her sonuç orijinal index'yi taşır. Aşağı akış işlemi için sıra önemliyse, arabellek sıralı bir yapıya neden olur ve bitişik çalıştırmaları temizler (TCP paketinin yeniden birleştirilmesi gibi).
Akışı kontrol noktası oluşturmayla birleştirebilir miyim?
Evet. Her sonucu geldiğinde bir kontrol noktası dosyasına ekleyin. Devam ederken denetim noktasını yükleyin, tamamlanan dizinleri filtreleyin ve yalnızca kalan görevleri yeniden işleyin.
İlgili Makaleler
Sonraki Adımlar
CAPTCHA çözümlerini geldikleri anda işleyin;CaptchaAI API anahtarınızı alınve akış boru hatları oluşturun.
İlgili kılavuzlar: