Kumpulan 500 tugas CAPTCHA diselesaikan secara tidak merata – ada yang diselesaikan dalam 8 detik, ada yang membutuhkan waktu 45. Menunggu setiap tugas selesai sebelum memproses hasil akan membuang-buang waktu antara solusi pertama dan terakhir. Streaming memungkinkan pipeline hilir Anda mengonsumsi setiap hasil saat solusi tersebut tiba.
Streaming vs. Proses Batch-Lalu
| Pendekatan | Saatnya untuk Hasil Pertama | Memori | Latensi Saluran Pipa |
|---|---|---|---|
| Tunggu semuanya | Setelah tugas paling lambat | Semua hasil dalam memori | Tinggi |
| Streaming sebagai terpecahkan | Setelah tugas respons kompetitif | Satu hasil pada satu waktu | Rendah |
| Micro-batch (potongan 10) | Setelah potongan pertama | 10 hasil sekaligus | Sedang |
Python: Generator Async untuk Hasil Streaming
Menggunakan asyncio dan aiohttp, setiap solusi segera dihasilkan melalui generator async:
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())
Instal dependensi:
pip install aiohttp
JavaScript: Pola Streaming EventEmitter
Node.js menggunakan pendekatan event-driven - memancarkan setiap hasil saat diselesaikan:
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);
Kapan Menggunakan Streaming vs. Kumpulkan-Semua
| Skenario | Pendekatan |
|---|---|
| Pengiriman formulir menggunakan token | Streaming – kirimkan setiap formulir segera setelah token tiba |
| Ekspor CSV dari semua hasil | Kumpulkan semua – tulis satu kali setelah batch selesai |
| Dashboard dengan kemajuan langsung | Streaming - perbarui UI pada setiap peristiwa hasil |
| Batch dengan ketergantungan antar tugas | Kumpulkan semua — proses secara berurutan setelah selesai |
| Batch besar (1.000+) | Streaming – mengurangi penggunaan memori puncak |
Pemecahan Masalah
| Masalah | Penyebab | Solusi |
|---|---|---|
| Hasil tiba dalam urutan acak | Normal – streaming menghasilkan yang respons kompetitif terlebih dahulu | Gunakan result.index untuk memetakan kembali ke tugas awal |
| Memori masih bertambah selama streaming | Menyimpan semua hasil dalam array | Memproses dan membuang hasilnya di handler |
| Hasil pertama memakan waktu terlalu lama | Semua tugas diserahkan secara bersamaan | Pengiriman terhuyung-huyung dengan batas semaphore atau konkurensi |
| Peringatan EventEmitter: MaxListenersExceeded | Terlalu banyak pendengar di streaming | Gunakan setMaxListeners() atau pastikan satu pendengar per jenis acara |
| Generator asinkron hang | Tugas yang belum terselesaikan dalam kumpulan yang tertunda | Tambahkan timeout ke poll_task; memastikan semua masa depan lengkap atau kesalahan |
Pertanyaan Umum
Apakah streaming meningkatkan panggilan API dibandingkan batch?
Tidak - jumlah panggilan kirim dan polling yang sama terjadi. Streaming hanya berubah saat aplikasi Anda memproses setiap hasil, bukan berapa banyak panggilan API yang dilakukan.
Bagaimana cara menjaga urutan tugas saat streaming?
Setiap hasil membawa index aslinya. Jika urutan penting untuk pemrosesan hilir, buffer menghasilkan struktur yang diurutkan dan menjalankan proses yang berdekatan (seperti perakitan ulang paket TCP).
Bisakah saya menggabungkan streaming dengan pos pemeriksaan?
Ya. Tambahkan setiap hasil ke file pos pemeriksaan saat hasilnya tiba. Saat melanjutkan, muat pos pemeriksaan, filter indeks yang sudah selesai, dan proses ulang hanya tugas yang tersisa.
Artikel Terkait
- Kafka Captchaai Streaming Pemrosesan Captcha
- Pemrosesan Prioritas Captcha Batch Antrian
Langkah Selanjutnya
Proses solusi CAPTCHA saat solusi tersebut tiba —dapatkan kunci API CaptchaAI Andadan membangun jaringan pipa streaming.
Panduan terkait: