DevOps & Scaling

Apache Kafka + CaptchaAI: Streaming Pemrosesan Task CAPTCHA

Begini pola yang bertahan ketika volume solve CAPTCHA naik dari ratusan menjadi puluhan ribu task per hari: pisahkan pengiriman task dari eksekusi solve lewat dua topik Kafka, biarkan consumer group menambah worker sendiri saat traffic naik, dan simpan history pesan supaya task yang gagal bisa diputar ulang tanpa kehilangan data. Tim otomatisasi dan agensi price-monitoring di Indonesia yang menjalankan scraper untuk ratusan target sekaligus biasanya baru merasakan titik ini — Redis atau RabbitMQ masih cukup untuk ribuan task/jam, tapi begitu volume naik ke puluhan ribu dan Anda butuh replay pesan saat worker crash, Kafka jadi pilihan yang lebih masuk akal.

Panduan ini membangun pipeline lengkap: dua topik Kafka (captcha-tasks dan captcha-results), producer di sisi scraper, worker consumer yang memanggil API CaptchaAI, lalu cara scaling worker dan memantau consumer lag supaya pipeline tidak diam-diam macet.

Arsitektur: Pisahkan Pengiriman Task dari Proses Solve

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

Dua topik Kafka menjaga producer dan consumer tetap independen satu sama lain:

  • captcha-tasks — parameter CAPTCHA (sitekey, pageurl, method) yang menunggu di-solve
  • captcha-results — token hasil solve, siap dipakai downstream tanpa scraper perlu menunggu langsung

Worker CAPTCHA membaca dari topik pertama, memanggil CaptchaAI, lalu menulis hasilnya ke topik kedua. Scraper yang mengirim task tidak perlu tahu siapa yang memprosesnya, dan consumer hasil tidak perlu tahu siapa yang mengirim task — keduanya hanya berkomunikasi lewat Kafka.

Sebelum Mulai: Setup Broker dan Dependency

# Python
pip install kafka-python requests

# Node.js
npm install kafkajs axios

Broker Kafka jalan di localhost:9092 pada contoh ini — ganti dengan alamat cluster Anda kalau sudah production, baik self-managed di VPS maupun layanan terkelola seperti Confluent Cloud atau Amazon MSK.

Langkah 1: Bikin Topik captcha-tasks dan captcha-results

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

Enam partisi berarti hingga enam consumer bisa jalan paralel dalam satu grup — angka ini yang nanti menentukan plafon jumlah worker Anda di bagian scaling.

Langkah 2: Producer — Scraper Mengirim Task CAPTCHA

Producer berjalan di sisi scraper: begitu scraper menemukan CAPTCHA di halaman, ia langsung publish task ke topik captcha-tasks dan lanjut ke pekerjaan berikutnya tanpa menunggu hasil solve selesai.

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

Langkah 3: Worker — Consumer yang Memanggil API CaptchaAI

Worker menjalankan pola empat langkah yang sama di semua bahasa: kirim task ke in.php, simpan task_id (captcha_id) dari respons, polling res.php sampai statusnya berubah, lalu pakai token yang keluar. Consumer baru commit offset Kafka setelah hasil berhasil dipublish ke captcha-results — supaya worker yang crash di tengah jalan tidak diam-diam kehilangan task.

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

Scaling Worker Sesuai Volume Task

# 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

Consumer group Kafka membagi partisi secara otomatis: tambah worker, Kafka rebalance sendiri tanpa Anda perlu campur tangan. Plafonnya adalah jumlah partisi topik — di atas enam worker pada contoh ini, worker ketujuh menganggur sampai Anda menambah partisi.

Samakan angka ini dengan alokasi thread CaptchaAI Anda, karena CaptchaAI menagih per thread konkuren, bukan per solve — bukan model bayar-per-CAPTCHA. Enam worker paralel yang masing-masing memproses beberapa task sekaligus biasanya cocok dengan paket ADVANCE ($90/bulan, 50 thread), dengan headroom cukup sebelum Anda perlu naik ke PREMIUM ($170/bulan, 100 thread). Kalau worker Anda deploy di region seperti AWS ap-southeast-1 (Singapura) atau GCP asia-southeast2 (Jakarta), latensi jaringan ke CaptchaAI dan ke broker Kafka biasanya jadi bottleneck lebih dulu ketimbang jumlah thread yang Anda punya.

Pantau Consumer Lag Sebelum Jadi Masalah

kafka-consumer-groups.sh --describe --group captcha-workers \
  --bootstrap-server localhost:9092
Metrik Kondisi sehat Perlu perhatian
Consumer lag < 100 > 1000 → tambah worker
Pesan masuk/detik Mengikuti laju scraper Lonjakan tiba-tiba menandakan burst traffic
Pesan keluar/detik Mengikuti laju masuk Tertinggal dari laju masuk = bottleneck di worker

Masalah Umum dan Cara Mengatasinya

Masalah Penyebab Solusi
Consumer lag terus naik Worker tidak sanggup mengimbangi laju task masuk Tambah worker instance (maksimal sejumlah partisi)
Hasil solve duplikat Worker crash sebelum commit offset Tambahkan pengecekan idempotency pada task_id di consumer hasil
Rebalancing terlalu sering Worker crash/restart berulang Naikkan session.timeout.ms; cek kemungkinan OOM
Task tidak terbagi rata ke partisi Distribusi key kurang baik Pakai key acak atau tambah jumlah partisi

Pertanyaan Umum

Kapan saya butuh Kafka, bukan cukup Redis atau RabbitMQ?

Kalau volume Anda masih di bawah kira-kira 1.000 task/jam, Redis atau RabbitMQ biasanya sudah cukup dan setup-nya jauh lebih ringan. Kafka baru terasa gunanya begitu Anda butuh replay pesan (worker crash tidak boleh menghilangkan task), throughput ratusan ribu pesan/detik, atau beberapa consumer group berbeda yang membaca stream task yang sama untuk keperluan masing-masing (solve, logging, alerting).

Berapa partisi dan worker yang ideal untuk volume task saya?

Mulai dari jumlah partisi yang sedikit lebih besar dari worker yang Anda rencanakan — enam partisi untuk enam worker adalah titik awal yang wajar. Menambah partisi jauh lebih mudah daripada menguranginya, jadi lebih baik sedikit berlebih di awal daripada pas-pasan. Samakan juga dengan jumlah thread CaptchaAI yang Anda alokasikan, supaya worker tidak menunggu antrean solve di sisi API.

Bagaimana memantau consumer lag tanpa harus cek dashboard terus-menerus?

Jadwalkan kafka-consumer-groups.sh --describe lewat cron, parsing kolom lag-nya, lalu kirim ke bot Telegram atau webhook Slack tim Anda kalau angkanya lewat ambang batas — banyak tim otomatisasi di Indonesia sudah punya pola notifikasi cron semacam ini untuk keperluan lain, jadi tinggal disambungkan ke command yang sama. Kalau pipeline Anda sudah punya stack monitoring, ekspor metrik JMX Kafka ke Prometheus jadi opsi yang lebih rapi.

Bagaimana menangani CAPTCHA yang gagal terus-menerus?

Batasi jumlah retry di level worker. Setelah retry maksimum tercapai, publish task tersebut ke topik captcha-dead-letter untuk diperiksa manual — jangan biarkan retry tanpa batas memblokir partisi dan menahan task lain yang antre di belakangnya.

Apakah pipeline scraping seperti ini aman dari sisi kepatuhan data?

CaptchaAI menyelesaikan CAPTCHA di halaman yang Anda proses; kepatuhan data tetap tanggung jawab Anda sebagai operator pipeline. Untuk tim yang beroperasi di Indonesia, pastikan scraping hanya menyentuh data yang memang berwenang Anda proses, sejalan dengan UU Pelindungan Data Pribadi (UU 27/2022) dan UU ITE — terutama kalau hasil scraping menyentuh data pribadi milik pihak lain.


Artikel Terkait

Komentar dinonaktifkan untuk artikel ini.