DevOps & Scaling

Apache Kafka + CaptchaAI: Akış CAPTCHA Görev İşleme

Saatte 5.000'den fazla CAPTCHA görevi işleyen bir kazıma hattında tek bir kuyruk çabuk tıkanır: üretici tarafı darboğaza girer, tüketici geriden gelir ve tek bir worker çökmesi tüm hattı durdurur. Apache Kafka bu sorunu, görev üretimini sonuç işlemeden ayırarak çözer — dayanıklı, tekrar oynatılabilir ve yatay olarak ölçeklenen bir mesaj akışı üzerinden.

Örneğin İstanbul merkezli bir e-ticaret fiyat izleme ekibi, farklı kaynaklardan gelen binlerce reCAPTCHA ve Turnstile görevini tek bir Python worker'ıyla değil, partition'lara dağıtılmış bir worker havuzuyla işler; bir worker çökse bile Kafka grubu kalan partition'ları otomatik olarak yeniden dağıtır ve hat durmaz. Bu rehberde iki topic'lik bir mimari kuracak, Python ve Node.js ile üretici/worker kodunu yazacak ve CaptchaAI API'sine bağlayacaksınız.

Akış Mimarisi: İki Topic, Tek Sorumluluk

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

İki topic neden gerekli?

Bu mimaride iki Kafka topic'i sorumlulukları net şekilde ayırıyor:

  • captcha-tasks — Henüz çözülmemiş CAPTCHA parametrelerini taşır (sitekey, pageurl, method).
  • captcha-results — Çözülmüş token'ları, aşağı yöndeki servislerin kullanımına hazır şekilde tutar.

Üretici (producer) tarafı sonuç tüketicilerinin nasıl çalıştığını bilmek zorunda değildir; bu ayrım iki tarafı birbirinden bağımsız olarak ölçeklendirmenizi sağlar.

Gerekli Bağımlılıklar

# Python
pip install kafka-python requests

# Node.js
npm install kafkajs axios

İpucu: Yerel geliştirmede tek node'luk bir Kafka örneği yeterlidir. Üretime taşımadan önce en az 3 broker'lı bir cluster planlayın — tek broker'da replikasyon olmadığı için replication-factor 1 yalnızca test amaçlıdır.

Adım 1: Kafka Topic'lerini Oluşturun

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

Altı partition, worker grubu başına altı paralel tüketiciye kadar izin verir. Partition sayısını sonradan artırmak mümkündür, azaltmak mümkün değildir — bu yüzden ilk kurulumda biraz payla başlamak mantıklıdır.

Adım 2: Görev Üretici Servisi (Kazıyıcı Tarafı)

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

Üretici her görevi task_id anahtarıyla gönderir; böylece aynı görev her zaman aynı partition'a düşer ve mesaj sırası korunur. acks="all" ayarı mesajın tüm replikalara yazılmasını bekler — tek node'lu yerel testte etkisi sınırlıdır ama üretimde veri kaybını önler.

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

Node.js tarafında kafkajs istemcisi aynı sözleşmeyi izler: task_id anahtarı, method/sitekey/pageurl alanları. Üretici ve worker'ın farklı dillerde yazılmış olması sorun değildir — JSON sözleşmesi sabit kaldığı sürece.

Üretim Ortamı Doğrulama Kuralları

  • Çözücü türü (method), sitekey/pageurl gibi hedef meta verileri veya sonuç yönlendirme bilgisi eksik olan mesajları, topic'e girmeden önce reddedin.
  • Aynı sorunun yeniden denemelerde tekrar tekrar çözülmeye çalışılmaması için üretici tarafında idempotency key kullanın (task_id bu amaca zaten hizmet eder).
  • Hatalı biçimlendirilmiş kayıtları, sonradan tekrar oynatabilmeniz için yeterli bağlamla birlikte captcha-dead-letter topic'ine yönlendirin.

Adım 3: CAPTCHA Worker'ı (Tüketici + Çözücü)

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

Worker, captcha-tasks'tan okur, CaptchaAI API'sine gönderir ve sonucu captcha-results'a yazar. enable_auto_commit=False kritik: offset yalnızca sonuç başarıyla yayınlandıktan sonra commit edilir, aksi halde bir çökme anında görev sessizce kaybolabilir.

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

Aynı worker mantığının Node.js karşılığı — method alanını görev JSON'ında değiştirerek CaptchaAI'nin desteklediği reCAPTCHA v2/v3, Cloudflare Turnstile ve GeeTest v3 gibi türleri aynı döngüyle işleyebilirsiniz.

Worker'ları Ölçeklendirme

# 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 tüketici grupları partition'ları worker'lar arasında otomatik dağıtır — elle bir atama yapmanıza gerek yoktur. Partition sayısının ötesine worker eklemenin faydası yoktur, çünkü boşta kalan worker'lar oluşur; bunun yerine partition sayısını artırın. Kural basit: paralel işleme hedefiniz neyse partition sayınız en az o kadar olmalı.

Kafka Tüketici Gecikmesini İzleme

kafka-consumer-groups.sh --describe --group captcha-workers \
  --bootstrap-server localhost:9092
Metrik Sağlıklı Uyarı
Tüketici gecikmesi < 100 > 1000 (worker ekleyin)
Gelen mesaj/sn Kazıyıcı hızıyla eşleşir Ani sıçramalar patlama (burst) olduğunu gösterir
Çıkan mesaj/sn Gelen hızla eşleşir Geride kalmak darboğaz demektir

Not: Consumer lag'i Prometheus + Kafka Exporter gibi bir izleme yığınına bağlarsanız eşik aşımında otomatik uyarı alırsınız; kafka-consumer-groups.sh çıktısını elle takip etmek zorunda kalmazsınız.

Sorun Giderme

Sorun Sebep Düzeltme
Tüketici gecikmesi artıyor Worker'lar görev hızına yetişemiyor Partition sayısına kadar daha fazla worker örneği ekleyin
Yinelenen sonuçlar Worker, offset'i commit etmeden çöküyor Sonuç tüketicisinde task_id üzerinden idempotency kontrolü ekleyin
Çok sık yeniden dengeleme Worker'lar çöküyor veya yeniden başlıyor session.timeout.ms değerini artırın; OOM olup olmadığını kontrol edin
Görevler eşit dağıtılmıyor Anahtar dağılımı kötü Rastgele anahtarlar veya daha fazla partition kullanın

Sık Sorulan Sorular

Kaç partition ve worker ile başlamalıyım?

Çoğu orta ölçekli kazıma hattı için 6 partition iyi bir başlangıçtır — 6 paralel worker'a kadar ölçeklenebilir. Saatlik hacminiz 10.000 görevi geçiyorsa partition sayısını 12–24 arasına çıkarın; sayıyı sonradan artırabilirsiniz, azaltamazsınız.

Bu ölçekte CaptchaAI için hangi plan yeterli?

Thread sayınız eşzamanlı worker sayınıza bağlıdır, partition sayısına değil. Saatte birkaç bin görev işleyen bir kurulum genellikle ADVANCE ($90/ay, 50 thread) ile rahat çalışır; 6 worker'lık küçük bir grup için STANDARD ($30/ay, 15 thread) bile başlangıç için yeterli olabilir. Hacim büyüdükçe üst kademelere geçebilirsiniz.

Redis veya RabbitMQ yerine neden Kafka?

Mesaj dayanıklılığı (replay), yüksek verim (saniyede 100.000+ mesaj) ve tüketici grubu ölçeklendirmesi gerektiğinde Kafka doğru seçimdir. Saatte 1.000 görevin altındaki daha basit kurulumlarda Redis veya RabbitMQ yeterli olur ve operasyonel yükü daha düşüktür.

Kafka mesajlarını ne kadar süre saklamalıyım (retention)?

captcha-tasks için birkaç saatlik bir retention genellikle yeterlidir — görev zaten worker tarafından hızla tüketilir. captcha-results topic'inde retention'ı, aşağı yöndeki servisin en yavaş tüketim hızına göre en az 24–48 saat tutmak, hata ayıklama ve tekrar oynatma için pay bırakır.

Üretici ve worker farklı dillerde yazılabilir mi?

Evet. Yukarıdaki örneklerde üretici Python, worker Node.js olabilir ya da tam tersi — Kafka mesaj formatı JSON olduğu sürece dil karışımı sorun yaratmaz. Ekipler genelde kazıyıcının zaten yazıldığı dili üretici için, operasyon ekibinin alışkın olduğu dili worker için seçer.

İlgili Makaleler

  • Toplu CAPTCHA sonuçlarını akış halinde işleme

Sonraki Adımlar

Akış CAPTCHA hattınızı bugün kurun — CaptchaAI API anahtarınızı alın ve yüksek verimli işleme için Kafka'ya bağlayın.

Benzer mimari rehberler:

Bu makale için yorumlar devre dışı bırakılmıştır.