DevOps & Scaling

Apache Kafka + CaptchaAI: потоковая обработка задач CAPTCHA

Очередь в памяти держится, пока задач не больше пары сотен в час; на тысячах начинаются потери при перезапуске воркера. Apache Kafka снимает это: сообщения хранятся на диске, порядок внутри партиции сохраняется даже при перезапуске воркера. Ниже — архитектура на два топика, код на Python и Node.js, и как считать число воркеров вместе с тарифом.

Архитектура: два топика вместо одной очереди

Продюсер и воркеры общаются только через Kafka:

[Scrapers] → Produce → [Kafka: captcha-tasks topic]
                              ↓
                    [CAPTCHA Worker Group]
                    (consume tasks, solve via CaptchaAI)
                              ↓
                    Produce → [Kafka: captcha-results topic]
                              ↓
                    [Result Consumer Group]
                    (process solutions, update database)
Топик Что в нём лежит
captcha-tasks параметры CAPTCHA, отправленные парсером
captcha-results решённые токены, готовые к записи в БД

Когда вообще нужен Kafka

Очередь — не самоцель, выбор инструмента зависит от объёма:

Задач в час Инструмент Почему
до ~1000 Redis или RabbitMQ Kafka здесь избыточна — хватает простой очереди без партиций
1000–20 000 Kafka, 5–15 воркеров нужен реплей сообщений и строгий порядок внутри партиции
20 000+ Kafka, отдельные топики по типу CAPTCHA или ICP один топик с шестью партициями перестаёт масштабироваться линейно

Что понадобится перед началом

Три вещи перед стартом:

  1. брокер Kafka — локальный localhost:9092 или адрес вашего кластера;
  2. клиентские библиотеки под ваш язык;
  3. рабочий API-ключ CaptchaAI.
# Python
pip install kafka-python requests

# Node.js
npm install kafkajs axios

Шаг 1. Создайте топики

kafka-topics.sh --create --topic captcha-tasks \
  --partitions 6 --replication-factor 1 \
  --bootstrap-server localhost:9092

kafka-topics.sh --create --topic captcha-results \
  --partitions 6 --replication-factor 1 \
  --bootstrap-server localhost:9092

Шесть партиций — до шести параллельных воркеров в группе (подробнее о выборе числа партиций — ниже).

Шаг 2. Продюсер задач на стороне парсера

Python

import json
from kafka import KafkaProducer

producer = KafkaProducer(
    bootstrap_servers=["localhost:9092"],
    value_serializer=lambda v: json.dumps(v).encode("utf-8"),
    key_serializer=lambda k: k.encode("utf-8") if k else None,
    acks="all",  # Wait for all replicas to confirm
    retries=3
)


def enqueue_captcha(task_id, sitekey, pageurl, captcha_type="userrecaptcha"):
    """Send a CAPTCHA task to Kafka."""
    task = {
        "task_id": task_id,
        "method": captcha_type,
        "sitekey": sitekey,
        "pageurl": pageurl,
        "submitted_at": __import__("time").time()
    }

    future = producer.send(
        "captcha-tasks",
        key=task_id,  # Key ensures same task goes to same partition
        value=task
    )
    future.get(timeout=10)  # Block until confirmed
    return task_id


# Submit tasks
enqueue_captcha("task_001", "6Le-wvkSAAAAAPBMRTvw0Q4Muexq9bi0DJwx_mJ-", "https://example.com")
enqueue_captcha("task_002", "6Le-wvkSAAAAAPBMRTvw0Q4Muexq9bi0DJwx_mJ-", "https://example.com")
producer.flush()

JavaScript

const { Kafka } = require("kafkajs");

const kafka = new Kafka({
  clientId: "captcha-producer",
  brokers: ["localhost:9092"],
});

const producer = kafka.producer();

async function enqueueCaptcha(taskId, sitekey, pageurl) {
  await producer.connect();

  const task = {
    task_id: taskId,
    method: "userrecaptcha",
    sitekey: sitekey,
    pageurl: pageurl,
    submitted_at: Date.now(),
  };

  await producer.send({
    topic: "captcha-tasks",
    messages: [{ key: taskId, value: JSON.stringify(task) }],
  });
}

(async () => {
  await enqueueCaptcha(
    "task_001",
    "6Le-wvkSAAAAAPBMRTvw0Q4Muexq9bi0DJwx_mJ-",
    "https://example.com"
  );
  await producer.disconnect();
})();

Ключ task_id направляет одинаковые задачи в одну партицию — повторная отправка не разъедет обработку между воркерами.

Валидация сообщений на стороне продюсера

  • отклоняйте сообщения без method, sitekey или pageurl;
  • добавляйте ключ идемпотентности, чтобы повторы не плодили дублей;
  • некорректные записи шлите в dead-letter топик.

Партиции Kafka и потоки CaptchaAI: как не упереться в лимит

Пока воркер ждёт ответа от res.php, он занимает поток CaptchaAI — второй потолок помимо партиций. План ADVANCE ($90/мес, 50 потоков) держит до 50 воркеров параллельно. Воркеров больше, чем потоков, — отставание растёт не из-за Kafka. Для команд с долларовыми счетами стоимость по потокам удобнее оплаты за решение.

Пример расчёта: при 5000 задачах в час и среднем времени решения 8 с одному воркеру хватает пропускной способности примерно на 450 задач в час, значит нужно 10–12 воркеров с запасом на пиковую нагрузку. Партиций стоит закладывать больше, чем этих 10–12 — иначе при следующем скачке объёма добавить воркеров без пересоздания топика не получится.

Если в pageurl или в теле задачи попадают реальные пользовательские данные (например, URL с параметрами сессии), отправляйте в Kafka только то, что вы вправе обрабатывать по 152-ФЗ «О персональных данных» (для читателей из РФ) или по аналогичным требованиям в вашей юрисдикции. Это не юридическая консультация, а базовая гигиена логирования при потоковой обработке.

Шаг 3. CAPTCHA-воркер: потребитель + решение через CaptchaAI

Python

import json
import os
import time
import requests
from kafka import KafkaConsumer, KafkaProducer

API_KEY = os.environ["CAPTCHAAI_API_KEY"]

consumer = KafkaConsumer(
    "captcha-tasks",
    bootstrap_servers=["localhost:9092"],
    group_id="captcha-workers",
    value_deserializer=lambda m: json.loads(m.decode("utf-8")),
    auto_offset_reset="earliest",
    enable_auto_commit=False,  # Manual commit after processing
    max_poll_records=10
)

result_producer = KafkaProducer(
    bootstrap_servers=["localhost:9092"],
    value_serializer=lambda v: json.dumps(v).encode("utf-8")
)


def solve_captcha(task):
    """Submit to CaptchaAI and poll for result."""
    # Submit
    resp = requests.post("https://ocr.captchaai.com/in.php", data={
        "key": API_KEY,
        "method": task["method"],
        "googlekey": task["sitekey"],
        "pageurl": task["pageurl"],
        "json": 1
    })
    data = resp.json()

    if data.get("status") != 1:
        return {"error": data.get("request")}

    captcha_id = data["request"]

    # Poll for result
    for _ in range(60):
        time.sleep(5)
        result = requests.get("https://ocr.captchaai.com/res.php", params={
            "key": API_KEY,
            "action": "get",
            "id": captcha_id,
            "json": 1
        }).json()

        if result.get("status") == 1:
            return {"solution": result["request"]}
        if result.get("request") != "CAPCHA_NOT_READY":
            return {"error": result.get("request")}

    return {"error": "TIMEOUT"}


# Main consumer loop
print("CAPTCHA worker started. Waiting for tasks...")
for message in consumer:
    task = message.value
    print(f"Processing {task['task_id']}...")

    result = solve_captcha(task)
    result["task_id"] = task["task_id"]
    result["solved_at"] = time.time()

    # Publish result
    result_producer.send("captcha-results", value=result)
    result_producer.flush()

    # Commit offset after successful processing
    consumer.commit()
    print(f"  → {task['task_id']}: {'solved' if 'solution' in result else result.get('error')}")

JavaScript

const { Kafka } = require("kafkajs");
const axios = require("axios");

const API_KEY = process.env.CAPTCHAAI_API_KEY;

const kafka = new Kafka({
  clientId: "captcha-worker",
  brokers: ["localhost:9092"],
});

const consumer = kafka.consumer({ groupId: "captcha-workers" });
const producer = kafka.producer();

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

async function solveCaptcha(task) {
  const submitResp = await axios.post(
    "https://ocr.captchaai.com/in.php",
    null,
    {
      params: {
        key: API_KEY,
        method: task.method,
        googlekey: task.sitekey,
        pageurl: task.pageurl,
        json: 1,
      },
    }
  );

  if (submitResp.data.status !== 1) {
    return { error: submitResp.data.request };
  }

  const captchaId = submitResp.data.request;

  for (let i = 0; i < 60; i++) {
    await sleep(5000);
    const result = await axios.get("https://ocr.captchaai.com/res.php", {
      params: { key: API_KEY, action: "get", id: captchaId, json: 1 },
    });

    if (result.data.status === 1) return { solution: result.data.request };
    if (result.data.request !== "CAPCHA_NOT_READY")
      return { error: result.data.request };
  }

  return { error: "TIMEOUT" };
}

async function run() {
  await consumer.connect();
  await producer.connect();
  await consumer.subscribe({ topic: "captcha-tasks", fromBeginning: false });

  await consumer.run({
    eachMessage: async ({ message }) => {
      const task = JSON.parse(message.value.toString());
      console.log(`Processing ${task.task_id}...`);

      const result = await solveCaptcha(task);
      result.task_id = task.task_id;
      result.solved_at = Date.now();

      await producer.send({
        topic: "captcha-results",
        messages: [{ value: JSON.stringify(result) }],
      });

      console.log(
        `  → ${task.task_id}: ${result.solution ? "solved" : result.error}`
      );
    },
  });
}

run();

Offset подтверждается только после публикации результата — упадёт воркер раньше, Kafka отдаст задачу другому.

Масштабирование воркеров

Consumer-группа сама распределяет партиции:

# 6 partitions, 3 workers → each worker gets 2 partitions
Worker-1: partitions 0, 1
Worker-2: partitions 2, 3
Worker-3: partitions 4, 5

# Add Worker-4 → rebalance
Worker-1: partitions 0, 1
Worker-2: partitions 2
Worker-3: partitions 3, 4
Worker-4: partition 5

Масштабируйтесь до числа партиций; сверх лимита воркер простаивает — сначала добавьте партиции.

Коротко:

  • партиций закладывайте с запасом относительно пиковых воркеров;
  • воркер сверх числа партиций простаивает — это ограничение самой Kafka, не CaptchaAI;
  • ребаланс при добавлении воркера занимает секунды, но кратковременно останавливает потребление.

Мониторинг отставания и пропускной способности

Ключевая метрика — отставание потребителей (consumer lag):

kafka-consumer-groups.sh --describe --group captcha-workers \
  --bootstrap-server localhost:9092
Метрика Норма Тревога
Отставание < 100 > 1000 — добавьте воркеров
Сообщ/сек, вход По темпу парсера Скачки = всплеск
Сообщ/сек, выход Как вход Отставание = узкое место

Типичные проблемы и как их исправить

Большинство инцидентов сводится к четырём сценариям:

Проблема Причина Исправление
Растёт отставание Воркеры не успевают Больше воркеров
Дубли результатов Воркер упал до offset Идемпотентность по task_id
Частый ребаланс Воркеры падают/рестартуют session.timeout.ms выше; проверьте OOM
Неравномерное распределение Плохие ключи Случайные ключи, больше партиций

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

Когда Kafka оправдана, а когда хватит Redis или RabbitMQ?

До ~1000 задач в час хватает Redis/RabbitMQ — проще в эксплуатации и не требует отдельного кластера. Kafka оправдана, когда нужен реплей сообщений (переиграть задачи с определённого момента) или когда воркеров уже десятки и нужна честная балансировка нагрузки между ними.

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

Больше, чем воркеров в пике нагрузки: для captcha-tasks закладывайте сразу 8–12 партиций, а не шесть. Увеличить число партиций у существующего топика можно, но перераспределение затронет ключи — проще заложить запас заранее.

Сколько потоков CaptchaAI нужно под поток задач?

При 5–15 с на решение одной CAPTCHA и 1000 задачах в час обычно хватает 10–15 воркеров, а значит и 10–15 потоков — план ADVANCE (50 потоков) даёт запас на пики и рост объёма без апгрейда тарифа.

Что делать при падении воркера — незавершённых CAPTCHA и дублях результатов?

Две независимые проблемы решаются двумя правилами:

  • CAPTCHA не решается после нескольких попыток — задайте лимит повторов на воркере (например, три попытки), а после исчерпания лимита публикуйте задачу в captcha-dead-letter для ручной проверки, не держите партицию заблокированной бесконечными ретраями;
  • воркер падает между публикацией результата и коммитом offset — публикуйте результат в captcha-results раньше, чем коммитите offset, и проверяйте task_id на стороне потребителя результатов, чтобы отфильтровать повторную обработку той же задачи.

Куда двигаться дальше

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