Tutorials

Потоковая передача результатов пакета: обработка решений CAPTCHA по мере их поступления

Токен, который уже решён, не обязан ждать соседей по пакету. Если конвейер забирает каждый результат в момент готовности, первая форма уходит на отправку через несколько секунд после старта, а не через минуту после того, как досчиталась самая медленная задача. Ниже — рабочая схема на Python и Node.js плюс правило, по которому выбирают между потоком и сбором всего пакета.

Разброс здесь не теоретический: в пакете из 500 задач часть решается за считаные секунды, часть — заметно дольше, в зависимости от типа проверки, нагрузки и сети. Ориентиры по времени решения берите из документации CaptchaAI, а архитектуру стройте так, чтобы медленный хвост не блокировал быстрое большинство.

Три стратегии: дождаться всех, стрим, микропакет

Подход Время до первого результата Память Задержка конвейера
Дождаться всех После самой медленной задачи Все результаты в памяти Высокая
Стрим по готовности После самой быстрой задачи Один результат за раз Низкая
Микропакет (куски по 10) После первого куска 10 результатов одновременно Средняя

Микропакет — компромисс. Он полезен, когда нижестоящая система принимает данные только батчами (запись в БД, публикация в очередь), но держать в памяти весь прогон вы не хотите.

Сколько потоков заказывать под такой конвейер

Потоковая отдача результатов не меняет тарификацию: CaptchaAI считает параллельные потоки, а не отдельные решения, и число решений на поток в тарифе не ограничено. Меняется только момент, когда ваш код видит токен.

Практический ориентир: лимит параллелизма в коде (max_concurrent, maxConcurrent) не должен превышать число потоков в тарифе — лишние задачи просто встанут в ожидание на стороне API.

  • BASIC ($15/мес, 5 потоков) — отладка схемы и небольшие прогоны.
  • ADVANCE ($90/мес, 50 потоков) — типичный парсинг-конвейер небольшой команды.
  • ENTERPRISE ($300/мес, 200 потоков) — постоянный поток задач в проде.

Ставьте лимит равным купленному числу потоков или чуть ниже — так остаётся запас на повторные попытки.

Python: асинхронный генератор, отдающий решения по одному

asyncio вместе с aiohttp даёт самый прямой вариант: генератор выдаёт результат, как только соответствующая корутина завершилась. Порядок выдачи — это порядок завершения, а не порядок отправки, поэтому индекс исходной задачи едет вместе с результатом.

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

Установите зависимость:

pip install aiohttp

Ключевая деталь — asyncio.wait(..., return_when=FIRST_COMPLETED): цикл просыпается на каждой завершённой задаче, а не на всём наборе. Обработчик внутри async for должен быть коротким. Тяжёлую работу — запись в БД, отправку формы — выносите в отдельную очередь, иначе медленный обработчик снова превратит поток в пакет.

JavaScript: тот же поток через EventEmitter

В Node.js естественная форма этого паттерна — событийная: класс складывает задачи в очередь, держит фиксированный параллелизм и эмитит событие result на каждом завершении.

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

Подписчик на result получает и успехи, и ошибки — счётчики и логирование живут в одном месте, а событие done закрывает прогон.

Когда поток не нужен

Сценарий Подход
Отправка форм по токену Поток — отправляйте форму сразу, как пришёл токен
Выгрузка всех результатов в CSV Сбор всего — записать один раз по завершении пакета
Панель с живым прогрессом Поток — обновлять интерфейс на каждом событии
Задачи с зависимостями между собой Сбор всего — обработать по порядку после завершения
Большие пакеты (1000+ задач) Поток — ниже пиковое потребление памяти

Правило простое: если следующий шаг конвейера умеет работать с одной записью — включайте поток. Если ему нужен целостный набор (сортировка, дедупликация, сводный отчёт) — собирайте всё.

Локальный сценарий: ночной прогон на нестабильном канале

Типичная ситуация для команды, которая гоняет сбор данных из европейского или центральноазиатского региона: RTT до целевого сайта плавает, часть задач уходит в повтор, а окно на прогон — несколько ночных часов.

При сборе всего пакета одна зависшая задача держит прогон целиком и съедает окно. При потоковой отдаче основная масса результатов уже записана и отправлена дальше, пока хвост из нескольких задач дорешивается или уходит в тайм-аут по poll_task. Второй плюс — контрольная точка: каждый результат сразу пишется в файл прогресса, и после обрыва канала прогон возобновляется с середины, а не с нуля.

Отдельная оговорка по данным: собирайте только то, что вы вправе обрабатывать. Для читателей из РФ ориентир — 152-ФЗ «О персональных данных», для трансграничных проектов — требования уровня GDPR. Это зона вашей юридической проверки, а не свойство сервиса решения CAPTCHA.

Поиск неисправностей

Проблема Причина Что делать
Результаты приходят вперемешку Так и задумано: поток отдаёт быстрые первыми Сопоставляйте по result.index с исходной задачей
Память всё равно растёт Все результаты складываются в массив Обрабатывайте и отбрасывайте результат прямо в обработчике
Первый результат идёт слишком долго Все задачи отправлены одновременно Ограничьте параллелизм семафором или размером окна
Предупреждение MaxListenersExceeded Слишком много подписчиков на потоке Один подписчик на тип события либо setMaxListeners()
Асинхронный генератор зависает В наборе осталась незавершённая задача Тайм-аут в poll_task; каждая корутина обязана завершиться результатом или ошибкой
Подряд идут ответы CAPCHA_NOT_READY Опрос чаще, чем сервис успевает решить Оставьте интервал опроса в 5 с и не уменьшайте его

Частые вопросы

Поток тратит больше запросов к API, чем обычный пакет?

Нет. Число вызовов in.php и res.php одинаково — меняется только момент, когда приложение обрабатывает готовый результат.

Как сохранить исходный порядок задач?

Каждый результат несёт свой index. Если порядок важен, буферизуйте результаты в отсортированной структуре и сбрасывайте вниз непрерывные участки — так же, как собирают TCP-сегменты.

Что делать, если задача зависла и не отдаёт результат?

Полагайтесь на тайм-аут, а не на бесконечный опрос: в примере на Python это параметр timeout=300 в poll_task, в Node.js — счётчик из 60 итераций. Зависшая корутина обязана завершиться ошибкой TIMEOUT, иначе генератор не закроется и прогон повиснет на одной задаче.

Какой лимит параллелизма ставить в коде?

Тот, что соответствует числу потоков в вашем тарифе. Значение max_concurrent=15 в примерах — иллюстрация: при 5 потоках BASIC ставьте 5, при 50 потоках ADVANCE — до 50.

Работает ли схема со всеми типами проверок?

Сама схема от типа не зависит — меняются только method и набор параметров при отправке. Применяйте её к поддерживаемым типам: reCAPTCHA v2 и v3, Cloudflare Turnstile и Cloudflare Challenge, GeeTest v3, image/OCR и grid. Для CaptchaFox (beta), Friendly Captcha (beta) и Lemin (beta) паттерн тот же, но статус — бета.


Что дальше

Комментарии для этой статьи отключены.