Eğitimler

Toplu Sonuçların Akışı: CAPTCHA Çözümlerini Geldiklerinde İşleme

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:

Bu makale için yorumlar devre dışı bırakılmıştır.