DevOps والتوسع

Apache Kafka + CaptchaAI: معالجة مهام CAPTCHA المتدفقة

الفكرة الجوهرية لمعالجة الكابتشا على نطاق واسع بسيطة: لا تدع الخدمة التي تجمع البيانات تنتظر حلّ كل اختبار. افصل إرسال المهمة عن استلام النتيجة، ودع Apache Kafka ينقل الرسائل بينهما. بهذا يواصل جامعو البيانات العمل بأقصى سرعة، بينما يتولّى عمّال مستقلّون الحلّ عبر CaptchaAI في الخلفية.

لنفترض أنك تدير منصّة لمقارنة أسعار الطيران في منطقة الخليج تسحب البيانات من عشرات المواقع المحميّة بـ reCAPTCHA. مع آلاف الطلبات في الساعة، تتحوّل قائمة الانتظار البسيطة إلى عنق الزجاجة. هنا تظهر قيمة Kafka: تقسيم إلى أقسام (partitioning)، ومجموعات مستهلكين (consumer groups)، وتتبّع للتأخر (lag)، وسجلّ رسائل قابل للإعادة. أما مع بضع مئات من المهام أو خدمة واحدة تُنتج وتستهلك، فقد تكفيك Redis أو قائمة أبسط.

متى يكون Kafka هو الخيار الصحيح؟

الواقع التشغيلي هل Kafka مناسب؟ السبب
بضع مئات من المهام يوميًا غالبًا لا تكلفة التشغيل والتعقيد أعلى من الفائدة
آلاف المهام في الساعة مع أكثر من مستهلك نعم يمنحك partitioning وconsumer groups طبيعيًا
تحتاج إعادة معالجة النتائج أو الاحتفاظ بسجل تدفق نعم Kafka يحافظ على log قابل للإعادة
خدمة واحدة صغيرة تريد queue سريعة فقط لا في الغالب قائمة أبسط تكفي

بنية خط المعالجة

يتكوّن الخط من ثلاث طبقات مستقلّة تتواصل فيما بينها عبر موضوعَي 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 في هذا المخطّط عبر أربع مراحل:

  1. يُنتج جامع البيانات مهمة إلى موضوع captcha-tasks ثم يعود إلى عمله فورًا دون انتظار.
  2. يسحب أحد العمّال المهمة من الموضوع، ويرسلها إلى CaptchaAI، وينتظر الرمز المحلول.
  3. يكتب العامل الرمز الناتج إلى موضوع captcha-results.
  4. تقرأ خدمة النتائج الرمز وتُكمل الخطوة النهائية، مثل تعبئة نموذج أو تحديث قاعدة بيانات.

نعتمد موضوعَين (topics) منفصلين لفصل المسؤوليات:

  • captcha-tasks - معلمات CAPTCHA في انتظار حلها
  • captcha-results - الرموز المميزة الجاهزة للاستخدام في المراحل النهائية

المتطلبات المسبقة

قبل تشغيل أول عامل، جهّز بيئتك بالمكوّنات التالية:

  • عنقود Kafka عامل، ولو وسيطًا واحدًا لأغراض التجربة المحلية.
  • مفتاح CaptchaAI API صالح مع رصيد Threads كافٍ لعدد العمّال المتوازين.
  • Python 3.9 أو أحدث، أو Node.js 18 أو أحدث، بحسب لغة التنفيذ.

ثم ثبّت المكتبات المطلوبة لكل بيئة:

# Python
pip install kafka-python requests

# Node.js
npm install kafkajs axios

شغّل وسيط Kafka على localhost:9092 (أو على عنوان مجموعتك).

الخطوة 1: إنشاء موضوعات Kafka

أنشئ موضوعَين منفصلين، أحدهما للمهام الواردة والآخر للنتائج المحلولة. عدد الأقسام الذي تختاره هنا هو ما يحدّد سقف التوازي لاحقًا، فاختره بحسب أقصى عدد عمّال تتوقّعه:

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: منتِج المهام (جهة جمع البيانات)

يتولّى المنتِج تحويل كل اختبار CAPTCHA يصادفه جامع البيانات إلى رسالة في موضوع captcha-tasks. المفتاح (key) هنا مهم: فهو يضمن وصول مهام المفتاح نفسه إلى القسم نفسه، ويحافظ على ترتيبها.

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

الخطوة 3: العامل الذي يستهلك المهمة ويحلّها

هذا هو قلب النظام: يقرأ العامل مهمة من captcha-tasks، ويرسلها إلى CaptchaAI عبر in.php، ثم يستطلع النتيجة من res.php حتى يجهز الرمز، وأخيرًا ينشره إلى captcha-results.

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')}")

الترتيب المهم هنا: نؤكّد الإزاحة (commit) فقط بعد نشر النتيجة إلى captcha-results، حتى لا تُفقد أي مهمة إذا توقّف العامل فجأة.

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

توسيع عدد العمّال

توزّع مجموعات مستهلكي Kafka الأقسام تلقائيًا على العمّال، وتعيد التوزيع (rebalance) كلما انضم عامل أو غادر:

# 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

يمكنك التوسّع حتى عدد الأقسام؛ ولزيادة التوازي أكثر، أضِف مزيدًا من الأقسام. ضع في حسبانك قاعدتين عمليتين:

  • لا فائدة من تشغيل عمّال أكثر من عدد الأقسام، إذ يبقى الفائض عاطلًا دون قسم يقرأ منه.
  • كل إضافة أو إزالة لعامل تُطلق إعادة توزيع قصيرة، لذا تجنّب تدوير العمّال بلا داعٍ أثناء ذروة الحِمل.

مراقبة الأداء عبر تأخر المستهلك

تأخر المستهلك (lag) هو أسرع مؤشر على تخلّف العمّال عن معدّل ورود المهام، وهو أول ما ينبغي أن تراقبه. تابع هذه الإشارات الثلاث معًا:

  • تأخر المستهلك المتراكم لكل قسم.
  • معدّل الرسائل الواردة إلى captcha-tasks.
  • معدّل الرسائل الصادرة إلى captcha-results.

اعرض تأخر مجموعة العمّال مباشرةً بالأمر التالي:

kafka-consumer-groups.sh --describe --group captcha-workers \
  --bootstrap-server localhost:9092
المقياس صحي تحذير
تأخر المستهلك (lag) < 100 > 1000 (أضِف عمّالاً)
الرسائل الواردة/الثانية تطابق معدل جمع البيانات الارتفاعات المفاجئة تدل على دفعة كبيرة
الرسائل الصادرة/الثانية تطابق معدل الورود التخلّف عن الركب = عنق زجاجة

حلّ المشكلات الشائعة

معظم مشكلات هذا الخط تعود إلى اختلال التوازن بين معدّل الورود وعدد العمّال، أو إلى توقيت تأكيد الإزاحة (commit). يلخّص الجدول التالي أكثرها شيوعًا وكيفية معالجتها:

المشكلة السبب الإجراء
lag يزداد باستمرار عدد العمّال أقل من حجم التدفق أو الاستطلاع بطيء جدًا زد عدد العمّال أو عدد الأقسام أو خفّض زمن الانتظار غير الضروري
المهام تُعاد معالجتها بعد إعادة التشغيل commit يحدث قبل اكتمال النشر إلى captcha-results أو بعده بشكل غير صحيح اجعل commit بعد نجاح معالجة الرسالة ونشر النتيجة
بعض العمّال عاطلون رغم وجود ضغط عدد العمّال أكبر من عدد الأقسام زد partitions أو خفّض عدد العمّال
رسائل النتائج تصل متأخرة رغم نجاح الحل يوجد عنق زجاجة في المنتج أو flush مفرط راقب زمن النشر إلى captcha-results ووازن بين flush وbatching

الأسئلة الشائعة

كم عدد الأقسام (partitions) التي أحتاجها لكل موضوع؟

ابدأ بعدد أقسام مساوٍ لأقصى عدد عمّال متوازين تتوقّعه، مع هامش للنمو. ستة أقسام تكفي حتى ستة عمّال؛ وإذا احتجت توازيًا أعلى، زد العدد (لا يمكن تقليصه). ولا تبالغ فيه، فكثرة الأقسام ترفع حمل التنسيق دون فائدة.

كيف أؤمّن مفتاح الـ API داخل العمّال الموزّعين؟

لا تضع مفتاح CaptchaAI داخل الكود أبدًا. مرّره عبر متغيّر بيئة مثل CAPTCHAAI_API_KEY، أو من خزنة أسرار (secrets manager) في الإنتاج. بهذا يقرأ كل عامل المفتاح نفسه، ويمكنك تدويره من مكان واحد دون إعادة النشر.

هل يحدّ عدد الـ Threads في خطة CaptchaAI من عدد العمّال المتوازين؟

نعم، وهو القيد الحقيقي على الإنتاجية. تُحاسب CaptchaAI على أساس عدد الـ Threads المتزامنة لا على عدد الحلول، مع حلول غير محدودة لكل Thread خلال الشهر. فخطة ADVANCE ($90 شهريًا، 50 Thread) تسمح بخمسين مهمة قيد الحل في اللحظة نفسها. أبقِ عدد المهام المتزامنة عبر عمّالك ضمن حصّة الـ Threads في خطتك.

ما أنواع الكابتشا التي يمكن للعامل حلّها في هذا المسار؟

يعالج العامل نفسه أي نوع تدعمه CaptchaAI بتغيير حقل method فقط: reCAPTCHA v2 وv3، وCloudflare Turnstile وChallenge، وGeeTest v3، والصور/OCR، والشبكات (grid)، وBLS، إضافةً إلى CaptchaFox وFriendly Captcha وLemin في مرحلة beta. أما hCaptcha وFunCaptcha (Arkose Labs) فغير مدعومَين حاليًا، وGeeTest v4 قيد الإعداد ولم يُطرح بعد.


الخطوات التالية

أدلة ذات صلة

التعليقات غير مفعّلة لهذا المقال.