Tutorial

Streaming Batch Result: Memproses Solve CAPTCHA Saat Tiba

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

Langkah Selanjutnya

Proses solusi CAPTCHA saat solusi tersebut tiba —dapatkan kunci API CaptchaAI Andadan membangun jaringan pipa streaming.

Panduan terkait:

Komentar dinonaktifkan untuk artikel ini.