Очередь в памяти держится, пока задач не больше пары сотен в час; на тысячах начинаются потери при перезапуске воркера. 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 | один топик с шестью партициями перестаёт масштабироваться линейно |
Что понадобится перед началом
Три вещи перед стартом:
- брокер Kafka — локальный
localhost:9092или адрес вашего кластера; - клиентские библиотеки под ваш язык;
- рабочий 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на стороне потребителя результатов, чтобы отфильтровать повторную обработку той же задачи.