Tutorials

Streaming Batch Result: Memproses Solve CAPTCHA Saat Tiba

Hasil solve pertama sebuah batch bisa Anda pakai jauh sebelum batch itu selesai — itu inti streaming. Pipeline Anda mengonsumsi setiap token pada detik token itu tersedia, bukan menunggu task terakhir kembali. Pada batch besar, jarak antara hasil pertama dan hasil terakhir sering menjadi porsi terbesar total waktu kerja.

Panduan ini membahas dua pola streaming yang lazim di produksi — async generator di Python dan EventEmitter di Node.js — plus cara menyetel concurrency sesuai jumlah thread Anda.

Kapan streaming layak dipakai, kapan sebaiknya tidak

Penentunya bukan ukuran batch, melainkan apa yang dilakukan proses hilir.

Skenario Pendekatan yang tepat
Token dipakai langsung untuk submit form Streaming — kirim form begitu token tiba
Ekspor CSV berisi seluruh hasil Kumpulkan semua — tulis satu kali saat batch selesai
Dashboard dengan progres langsung Streaming — perbarui UI pada setiap event hasil
Task saling bergantung satu sama lain Kumpulkan semua — proses berurutan setelah semuanya kembali
Batch besar (1.000+ task) Streaming — menekan puncak pemakaian memori

Bila dua baris berlaku sekaligus — token dipakai untuk submit form sekaligus diarsipkan ke CSV — streaming tetap pilihan yang benar: tulis file secara append di dalam handler, bukan sekali di akhir.

Tiga model pengambilan hasil batch

Model Hasil pertama tersedia Pemakaian memori Latensi pipeline
Tunggu semua task Setelah task paling lambat kembali Semua hasil ditahan di memori Tinggi
Streaming per hasil Setelah task pertama selesai Satu hasil pada satu waktu Rendah
Micro-batch (potongan 10) Setelah potongan pertama penuh 10 hasil sekaligus Sedang

Micro-batch cocok bila proses hilir lebih murah dijalankan secara bulk, misalnya INSERT database.

Streaming di Python dengan async generator

Dengan asyncio dan aiohttp, setiap hasil dikeluarkan lewat yield begitu task-nya selesai. Alurnya tetap empat langkah seperti solve tunggal — kirim, simpan task ID, polling, pakai token — hanya saja puluhan alur berjalan berdampingan dan hasil keluar sesuai urutan penyelesaian.

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())

Tiga bagian kode di atas yang menentukan perilakunya di produksi:

  • asyncio.Semaphore menjaga jumlah task bersamaan tetap sesuai jatah thread Anda.
  • asyncio.wait(..., FIRST_COMPLETED) membuat hasil keluar sesuai urutan penyelesaian, bukan urutan pengiriman.
  • timeout=300 pada poll_task mencegah satu task bermasalah menggantungkan seluruh generator.

Instal dependensinya:

pip install aiohttp

Streaming di Node.js dengan EventEmitter

Di Node.js polanya event-driven: stream memancarkan event result tiap task selesai dan done saat antrean habis, sehingga logika bisnis terpisah dari concurrency.

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);

Pemisahan itu membuat logika bisnis tidak perlu tahu apa pun soal concurrency:

  • result dipancarkan sekali per task, berisi index, status, dan token atau error.
  • done dipancarkan sekali saat completed sudah menyamai total — tempat yang tepat untuk menutup koneksi database atau menulis ringkasan.

Menyelaraskan concurrency dengan jumlah thread

Angka max_concurrent bukan angka bebas. CaptchaAI menagih per thread — satu thread berarti satu CAPTCHA yang sedang diproses — dengan solve tanpa batas per thread. Jadi batas concurrency ditentukan paket Anda, bukan besarnya server.

Paket Harga per bulan Thread max_concurrent yang wajar
BASIC $15 5 5
STANDARD $30 15 15
ADVANCE $90 50 45–50

Contoh yang lazim di pasar Indonesia: freelancer scraping dan tim price-monitoring kecil menjalankan worker di AWS ap-southeast-1 (Singapura) atau ap-southeast-3 (Jakarta), dengan beban menumpuk pada batch malam. Karena biaya tetap per thread, lonjakan malam tidak menambah tagihan; max_concurrent di atas jumlah thread hanya menambah antrean.

Tiga hal yang layak diperiksa sebelum menaikkan angkanya:

  • Sisakan satu atau dua thread bila API key yang sama juga dipakai worker lain atau cron malam.
  • Ukur dulu waktu penyelesaian rata-rata tipe CAPTCHA yang Anda pakai, lalu turunkan perkiraan throughput dari angka itu — bukan sebaliknya.
  • Naikkan concurrency bertahap dan pantau rasio task berstatus failed; lonjakan kegagalan biasanya berarti target Anda, bukan solver, yang menjadi hambatan.

Catatan kepatuhan lokal: di bawah UU Pelindungan Data Pribadi (UU 27/2022), proses hanya data yang berhak Anda olah.

Checklist sebelum menjalankan batch besar

  1. Uji pipeline dengan 5–10 task lebih dulu, baru naikkan volumenya.
  2. Pastikan max_concurrent tidak melebihi jumlah thread paket Anda.
  3. Simpan index bersama setiap hasil sejak awal — urutan keluaran berbeda dari urutan pengiriman.
  4. Tetapkan batas waktu pada setiap task, bukan hanya pada keseluruhan batch.
  5. Pakai token sesegera mungkin setelah tiba, terutama untuk reCAPTCHA yang masa berlakunya pendek.

Masalah umum saat streaming dan cara mengatasinya

Lima kendala berikut hampir selalu muncul saat pola streaming pertama kali masuk produksi.

Masalah Penyebab Solusi
Hasil datang dalam urutan acak Normal — streaming mengeluarkan hasil sesuai urutan penyelesaian Pakai result.index untuk memetakan kembali ke task asal
Memori tetap membengkak saat streaming Semua hasil masih ditumpuk dalam array Proses lalu buang hasilnya di dalam handler
Hasil pertama terlalu lama muncul Semua task dikirim serentak Sebar pengiriman dengan semaphore atau batas concurrency
Peringatan EventEmitter: MaxListenersExceeded Terlalu banyak listener pada satu stream Pakai setMaxListeners() atau pastikan satu listener per jenis event
Async generator menggantung Ada task yang tidak pernah selesai di set pending Tambahkan batas waktu pada poll_task; pastikan setiap future selesai atau error

Pertanyaan umum

Berapa nilai max_concurrent yang aman?

Samakan dengan jumlah thread paket Anda, sisakan ruang bila ada proses lain memakai API key sama. Menaikkannya lebih jauh tidak menambah throughput.

Apakah token dari streaming bisa langsung dipakai?

Justru itu keunggulannya. Token reCAPTCHA umumnya berlaku sekitar dua menit, jadi pola "kumpulkan dulu, pakai belakangan" berisiko membuat token kedaluwarsa.

Apakah pola ini berlaku untuk tipe CAPTCHA selain reCAPTCHA?

Ya. Alur kirim–polling–token identik untuk reCAPTCHA v2/v3, reCAPTCHA Enterprise, Cloudflare Turnstile, dan GeeTest v3; yang berubah hanya parameter method. Catatan: hCaptcha dan FunCaptcha belum didukung, GeeTest v4 masih segera hadir, sedangkan CaptchaFox, Friendly Captcha, dan Lemin berstatus beta.

Apa yang terjadi kalau satu task gagal total?

Tidak ada yang berhenti. Setiap task punya batas waktu sendiri dan mengembalikan objek berstatus failed, jadi stream tetap jalan. Catat index yang gagal dan jalankan ulang subset itu.

Berapa banyak task yang aman dikirim dalam satu batch?

Batas praktisnya bukan API, melainkan memori proses Anda dan masa berlaku token. Selama hasil diproses lalu dibuang di dalam handler, batch ribuan task berjalan wajar; yang perlu dijaga adalah jangan menumpuk token yang belum sempat dipakai.

Artikel Terkait

Langkah Selanjutnya

Pakai setiap hasil solve pada detik hasil itu tiba — ambil kunci API CaptchaAI Anda dan bangun pipeline streaming pertama Anda.

Panduan terkait:

Komentar dinonaktifkan untuk artikel ini.