Tutorials

Создание очереди решения CAPTCHA в Node.js

Единственный поток Node.js — не проблема для решения CAPTCHA: узкое место здесь не CPU, а ожидание ответа API, а с этим event loop справляется куда эффективнее классического многопоточного воркера.

Ниже — пять рабочих паттернов очереди для CaptchaAI: от Promise.allSettled на десяток задач до продакшен-очереди с приоритетами, повторами и метриками.


Какой паттерн очереди выбрать

Пять паттернов ниже расположены по нарастанию сложности — берите первый, который закрывает задачу, и переходите дальше только когда потребуется больше контроля:

  • Простой пакет (Promise.allSettled) — разовый скрипт или десяток задач за раз, без инфраструктуры.
  • Очередь с ограничением параллелизма — когда нужно удержаться в лимите потоков вашего тарифа.
  • Очередь на EventEmitter — когда важен прогресс в реальном времени: дашборд, лог, вебхук на фронтенд.
  • Приоритетная очередь — когда часть задач (например, checkout) не может ждать в общей очереди наравне с фоновым парсингом.
  • Очередь с повторами и dead-letter — для продакшена, где сетевые сбои и временные ошибки API неизбежны.

Быстрый старт: пакетное решение через Promise.allSettled

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

const API_KEY = "YOUR_API_KEY";

function sleep(ms) {
  return new Promise((r) => setTimeout(r, ms));
}

async function solveSingle(method, params) {
  const submitResp = await fetch("https://ocr.captchaai.com/in.php", {
    method: "POST",
    body: new URLSearchParams({ key: API_KEY, method, json: "1", ...params }),
  });
  const submitData = await submitResp.json();
  if (submitData.status !== 1) throw new Error(submitData.request);
  const taskId = submitData.request;

  for (let i = 0; i < 30; i++) {
    await sleep(5000);
    const pollResp = await fetch(
      `https://ocr.captchaai.com/res.php?${new URLSearchParams({
        key: API_KEY,
        action: "get",
        id: taskId,
        json: "1",
      })}`
    );
    const data = await pollResp.json();
    if (data.status === 1) return data.request;
    if (data.request === "ERROR_CAPTCHA_UNSOLVABLE") throw new Error("Unsolvable");
  }
  throw new Error("Timed out");
}

// Solve all at once
async function solveBatch(tasks) {
  const results = await Promise.allSettled(
    tasks.map((task) => solveSingle(task.method, task.params))
  );

  return results.map((result, i) => ({
    taskId: tasks[i].id,
    status: result.status,
    value: result.status === "fulfilled" ? result.value : null,
    error: result.status === "rejected" ? result.reason.message : null,
  }));
}

// Usage
const tasks = Array.from({ length: 10 }, (_, i) => ({
  id: i,
  method: "userrecaptcha",
  params: { googlekey: `KEY_${i}`, pageurl: `https://example.com/${i}` },
}));

const results = await solveBatch(tasks);
console.log(`Solved: ${results.filter((r) => r.status === "fulfilled").length}/10`);

Для десятка-другого задач этого достаточно: Promise.allSettled запускает всё сразу и не роняет весь пакет из-за одной ошибки. Проблема начинается на масштабе — без ограничения параллелизма вы почти сразу упрётесь в лимит одновременных вызовов API.


Ограничиваем параллелизм: своя очередь на классе

Плата CaptchaAI начисляется за одновременные потоки, а не за отдельное решение, поэтому значение maxConcurrent в очереди имеет смысл выставлять по тарифу: BASIC ($15/мес, 5 потоков) держит не больше пяти задач одновременно, ADVANCE ($90/мес, 50 потоков) — уже полсотни, а VIP-тарифы рассчитаны на промышленный парсинг с тысячами потоков.

Совет: если maxConcurrent в коде выше числа потоков в тарифе, часть вызовов начнёт получать ERROR_NO_SLOT_AVAILABLE вместо того, чтобы просто подождать своей очереди — рассинхрон между кодом и тарифом заметен сразу по логам.

class ConcurrencyQueue {
  constructor(maxConcurrent = 5) {
    this.maxConcurrent = maxConcurrent;
    this.running = 0;
    this.queue = [];
    this.results = [];
  }

  add(fn) {
    return new Promise((resolve, reject) => {
      this.queue.push({ fn, resolve, reject });
      this.#process();
    });
  }

  async #process() {
    if (this.running >= this.maxConcurrent || this.queue.length === 0) return;

    this.running++;
    const { fn, resolve, reject } = this.queue.shift();

    try {
      const result = await fn();
      resolve(result);
    } catch (error) {
      reject(error);
    } finally {
      this.running--;
      this.#process();
    }
  }

  async addBatch(fns) {
    return Promise.allSettled(fns.map((fn) => this.add(fn)));
  }
}

// Usage
const queue = new ConcurrencyQueue(5);

const tasks = Array.from({ length: 20 }, (_, i) => () =>
  solveSingle("userrecaptcha", {
    googlekey: `KEY_${i}`,
    pageurl: `https://example.com/${i}`,
  })
);

const results = await queue.addBatch(tasks);
const solved = results.filter((r) => r.status === "fulfilled");
console.log(`Solved: ${solved.length}/${results.length}`);

Очередь с событиями: прогресс в реальном времени через EventEmitter

Когда нужно не просто дождаться результата, а показывать прогресс по каждой задаче (дашборд, лог в консоль, вебхук на фронтенд), удобнее строить очередь поверх EventEmitter:

Событие complete срабатывает один раз — когда очередь пуста и нет ни одной активной задачи. Удобная точка, чтобы закрыть соединение с фронтендом или записать итоговую метрику.

const { EventEmitter } = require("events");

class CaptchaQueue extends EventEmitter {
  #apiKey;
  #maxConcurrent;
  #pending;
  #active;

  constructor(apiKey, maxConcurrent = 5) {
    super();
    this.#apiKey = apiKey;
    this.#maxConcurrent = maxConcurrent;
    this.#pending = [];
    this.#active = 0;
    this.stats = { submitted: 0, solved: 0, failed: 0 };
  }

  submit(id, method, params) {
    this.#pending.push({ id, method, params });
    this.stats.submitted++;
    this.emit("submitted", { id, total: this.stats.submitted });
    this.#drain();
  }

  async #drain() {
    while (this.#active < this.#maxConcurrent && this.#pending.length > 0) {
      const task = this.#pending.shift();
      this.#active++;
      this.#solve(task).finally(() => {
        this.#active--;
        this.#drain();
        if (this.#active === 0 && this.#pending.length === 0) {
          this.emit("complete", this.stats);
        }
      });
    }
  }

  async #solve(task) {
    try {
      const token = await solveSingle(task.method, task.params);
      this.stats.solved++;
      this.emit("solved", { id: task.id, token, stats: { ...this.stats } });
    } catch (error) {
      this.stats.failed++;
      this.emit("failed", { id: task.id, error: error.message, stats: { ...this.stats } });
    }
  }
}

// Usage
const queue = new CaptchaQueue("YOUR_API_KEY", 5);

queue.on("submitted", ({ id, total }) => {
  console.log(`Submitted #${id} (total: ${total})`);
});

queue.on("solved", ({ id, stats }) => {
  console.log(`Solved #${id} — ${stats.solved}/${stats.submitted}`);
});

queue.on("failed", ({ id, error }) => {
  console.log(`Failed #${id}: ${error}`);
});

queue.on("complete", (stats) => {
  const rate = ((stats.solved / stats.submitted) * 100).toFixed(1);
  console.log(`Done: ${stats.solved}/${stats.submitted} (${rate}%)`);
});

// Submit tasks
for (let i = 0; i < 15; i++) {
  queue.submit(i, "userrecaptcha", {
    googlekey: `KEY_${i}`,
    pageurl: `https://example.com/${i}`,
  });
}

Приоритеты: сначала checkout, потом парсинг

Не все задачи равнозначны: токен для оформления заказа блокирует пользователя прямо сейчас, а токен для фонового парсинга каталога может подождать. Приоритетная очередь решает это без отдельного микросервиса:

class PriorityQueue {
  #items = [];

  enqueue(item, priority) {
    this.#items.push({ item, priority });
    this.#items.sort((a, b) => a.priority - b.priority);
  }

  dequeue() {
    return this.#items.shift()?.item;
  }

  get length() {
    return this.#items.length;
  }
}

class PriorityCaptchaQueue {
  #apiKey;
  #maxConcurrent;
  #queue;
  #active;
  #results;

  constructor(apiKey, maxConcurrent = 5) {
    this.#apiKey = apiKey;
    this.#maxConcurrent = maxConcurrent;
    this.#queue = new PriorityQueue();
    this.#active = 0;
    this.#results = new Map();
  }

  submit(id, method, params, priority = 5) {
    return new Promise((resolve, reject) => {
      this.#queue.enqueue({ id, method, params, resolve, reject }, priority);
      this.#drain();
    });
  }

  async #drain() {
    while (this.#active < this.#maxConcurrent && this.#queue.length > 0) {
      const task = this.#queue.dequeue();
      this.#active++;

      solveSingle(task.method, task.params)
        .then((token) => {
          this.#results.set(task.id, { status: "solved", token });
          task.resolve(token);
        })
        .catch((err) => {
          this.#results.set(task.id, { status: "error", error: err.message });
          task.reject(err);
        })
        .finally(() => {
          this.#active--;
          this.#drain();
        });
    }
  }
}

// Usage: high-priority checkout, low-priority scraping
const pq = new PriorityCaptchaQueue("YOUR_API_KEY", 3);

// Priority 1 (highest) — checkout
const checkoutToken = pq.submit(
  "checkout_1",
  "turnstile",
  { sitekey: "KEY", pageurl: "https://shop.com/checkout" },
  1
);

// Priority 5 (normal) — product scraping
for (let i = 0; i < 5; i++) {
  pq.submit(
    `product_${i}`,
    "userrecaptcha",
    { googlekey: "KEY", pageurl: `https://shop.com/p/${i}` },
    5
  );
}

Повторные попытки и dead-letter очередь

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

  • Стоит повторить: таймаут сети, ERROR_NO_SLOT_AVAILABLE, разрыв соединения.
  • Не стоит повторять: ERROR_CAPTCHA_UNSOLVABLE и другие ошибки самой задачи — повтор с теми же параметрами даст тот же результат.
class RetryQueue {
  #apiKey;
  #maxRetries;
  #results;
  #deadLetter;

  constructor(apiKey, maxRetries = 3) {
    this.#apiKey = apiKey;
    this.#maxRetries = maxRetries;
    this.#results = [];
    this.#deadLetter = [];
  }

  async processBatch(tasks, maxConcurrent = 5) {
    const queue = tasks.map((t) => ({ ...t, attempts: 0 }));

    while (queue.length > 0) {
      const batch = queue.splice(0, maxConcurrent);
      const results = await Promise.allSettled(
        batch.map((task) => this.#solveWithRetry(task))
      );

      for (let i = 0; i < results.length; i++) {
        const result = results[i];
        const task = batch[i];

        if (result.status === "fulfilled") {
          this.#results.push({ id: task.id, token: result.value });
        } else {
          task.attempts++;
          if (task.attempts < this.#maxRetries) {
            queue.push(task); // Retry
            console.log(`Retry ${task.attempts}/${this.#maxRetries}: ${task.id}`);
          } else {
            this.#deadLetter.push({
              id: task.id,
              error: result.reason.message,
              attempts: task.attempts,
            });
          }
        }
      }
    }

    return {
      solved: this.#results,
      failed: this.#deadLetter,
    };
  }

  async #solveWithRetry(task) {
    return solveSingle(task.method, task.params);
  }
}

Метрики очереди: что стоит мониторить

Без метрик сложно понять, растёт ли очередь из-за пиковой нагрузки или из-за деградации самой интеграции. Минимальный набор — число решённых, число ошибок, среднее время решения и пропускная способность в минуту:

class QueueMonitor {
  #startTime;
  #solveTimes;

  constructor() {
    this.#startTime = Date.now();
    this.#solveTimes = [];
    this.counts = { submitted: 0, solving: 0, solved: 0, failed: 0 };
  }

  recordSubmit() {
    this.counts.submitted++;
    this.counts.solving++;
  }

  recordSolved(solveTime) {
    this.counts.solving--;
    this.counts.solved++;
    this.#solveTimes.push(solveTime);
  }

  recordFailed() {
    this.counts.solving--;
    this.counts.failed++;
  }

  report() {
    const elapsed = (Date.now() - this.#startTime) / 1000;
    const avgTime =
      this.#solveTimes.length > 0
        ? this.#solveTimes.reduce((a, b) => a + b, 0) / this.#solveTimes.length
        : 0;
    const throughput = this.counts.solved / (elapsed / 60);
    const successRate =
      this.counts.solved + this.counts.failed > 0
        ? (this.counts.solved / (this.counts.solved + this.counts.failed)) * 100
        : 0;

    return {
      elapsed: `${elapsed.toFixed(0)}s`,
      submitted: this.counts.submitted,
      solving: this.counts.solving,
      solved: this.counts.solved,
      failed: this.counts.failed,
      avgSolveTime: `${(avgTime / 1000).toFixed(1)}s`,
      throughput: `${throughput.toFixed(1)}/min`,
      successRate: `${successRate.toFixed(1)}%`,
    };
  }
}

Логируя задачи, сохраняйте только адрес страницы и токен — не персональные данные пользователя, чьё действие вызвало CAPTCHA. Для аудитории из РФ это требование 152-ФЗ, для проектов с европейскими пользователями — GDPR; правило одно: минимизируйте, что логируете.


Типичные проблемы и их решение

Большинство сбоев очереди сводится к пяти симптомам — ниже быстрая диагностика без разбора логов с нуля:

Симптом Причина Решение
Все promise разом переходят в rejected Упёрлись в лимит скорости API Понизьте maxConcurrent
Память растёт со временем Результаты накапливаются в массиве Периодически обрабатывайте и очищайте результаты
Очередь пуста, а задачи не завершены Не вызывается drain() после завершения Проверьте, что триггер слива стоит в блоке finally
ERROR_NO_SLOT_AVAILABLE Слишком много одновременных вызовов API Добавьте задержку между отправками
Dead-letter очередь растёт Ошибки повторяются раз за разом Проверьте тип ошибки — возможно, дело в параметрах запроса

Часто задаваемые вопросы

Сколько потоков тарифа нужно под maxConcurrent = 20?

Число потоков должно быть не меньше maxConcurrent, иначе часть задач будет ждать освобождения потока вместо решения. Под maxConcurrent = 20 подходит ADVANCE ($90/мес, 50 потоков) — с запасом на пиковую нагрузку.

Что означает ошибка ERROR_NO_SLOT_AVAILABLE?

Она приходит, когда одновременных запросов больше, чем разрешает тариф. Это не сбой очереди, а сигнал снизить maxConcurrent или перейти на тариф с большим числом потоков.

BullMQ или свой класс на JS — что выбрать?

Для одного процесса классов из этой статьи достаточно — без внешних зависимостей. BullMQ нужен, когда очередь должна переживать перезапуск процесса или её обслуживают несколько воркеров через общий Redis.

Как логировать очередь, не нарушая 152-ФЗ и GDPR?

Логируйте только то, что нужно для отладки: ID задачи, тип CAPTCHA, время решения, код ошибки. Персональные данные пользователя страницы с CAPTCHA для работы очереди не нужны — не сохраняйте их.

Меняется ли расчёт потоков очереди в зависимости от региона деплоя?

Число потоков в тарифе от региона не зависит, а вот RTT до ocr.captchaai.com — зависит: из европейских дата-центров и Казахстана опрос res.php обычно быстрее, чем при мобильном или нестабильном канале связи. На нестабильной сети закладывайте больший интервал sleep() между опросами и небольшой запас по потокам сверх расчётного maxConcurrent, чтобы всплеск повторов не упирался в лимит тарифа.


Коротко о главном

Node.js отлично справляется с параллельными операциями ввода-вывода — это ровно тот случай, для которого создаются очереди CAPTCHA с CaptchaAI. Начните с Promise.allSettled для простых пакетов, переходите на EventEmitter, когда нужен прогресс в реальном времени, добавляйте приоритеты для критичных для бизнеса потоков и очередь повторов — для устойчивости к временным сбоям API.

Смотрите также


Что читать дальше

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