DevOps và Mở Rộng

Apache Kafka + CaptchaAI: xử lý tác vụ CAPTCHA theo luồng

Redis chịu tốt vài trăm CAPTCHA mỗi giờ, nhưng qua vài nghìn tác vụ/giờ, ba vấn đề lộ ra: mất tác vụ khi worker crash, không replay được lịch sử, khó chia việc đều cho nhiều tiến trình song song. Kafka giải quyết đúng ba việc đó — độ bền tin nhắn, replay, mở rộng theo consumer group — nên đây là lựa chọn hạ tầng phổ biến khi ghép CaptchaAI vào pipeline dữ liệu quy mô lớn.

Ví dụ thực tế: một đội automation ở công ty outsourcing theo dõi giá hàng chục nghìn SKU trên Shopee, Lazada và Tiki mỗi ngày. Mỗi lần crawler gặp reCAPTCHA hoặc Turnstile, tác vụ được đẩy vào Kafka thay vì gọi CaptchaAI ngay trong luồng scraping — scraper không bị chặn chờ, còn nhóm worker riêng lo giải CAPTCHA độc lập và co giãn theo tải.

Kiến trúc: tách producer khỏi worker giải CAPTCHA

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

Hai topic Kafka tách bạch hai mối quan tâm khác nhau:

Topic Vai trò
captcha-tasks Tham số CAPTCHA đang chờ giải
captcha-results Token đã giải xong, sẵn sàng cho bước xử lý tiếp theo

Chuẩn bị môi trường

# Python
pip install kafka-python requests

# Node.js
npm install kafkajs axios

Kafka broker chạy tại localhost:9092 (hoặc địa chỉ cluster của bạn).

Bước 1: Tạo topic

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

Sáu partition cho phép tối đa sáu consumer chạy song song trong cùng một group — đây cũng là trần thông lượng của bước 3 bên dưới, trước khi bạn cần thêm partition.

Bước 2: Producer đẩy tác vụ CAPTCHA vào Kafka (phía scraper)

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

key: taskId đảm bảo mọi retry của cùng task_id rơi vào đúng một partition, giữ đúng thứ tự xử lý.

Nguyên tắc validate ở producer

  • Từ chối tác vụ thiếu method, pageurl hoặc thông tin định tuyến kết quả trước khi lọt vào topic.
  • Gắn key ổn định (task_id) để retry không sinh ra kết quả giải trùng.
  • Đẩy message sai định dạng sang topic dead-letter riêng, giữ ngữ cảnh để replay thủ công.

Bước 3: Worker CAPTCHA — vừa là consumer, vừa là solver

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

Lưu ý: consumer.commit() chỉ chạy sau khi kết quả đã publish thành công sang captcha-results — an toàn hơn commit sớm rồi mất kết quả nếu worker crash giữa chừng.

Mở rộng số lượng worker

Consumer group của Kafka tự động chia partition cho các worker đang chạy:

# 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

Hai điều cần nhớ khi scale:

  • Trần mở rộng đúng bằng số partition hiện có — thêm worker thứ 7 khi chỉ có 6 partition sẽ không nhận được việc gì.
  • Muốn scale tiếp thì tăng số partition trước, rồi mới thêm worker.

Giám sát pipeline qua consumer lag

Theo dõi độ trễ consumer là cách nhanh nhất để biết pipeline có đang tụt hậu:

kafka-consumer-groups.sh --describe --group captcha-workers \
  --bootstrap-server localhost:9092
Chỉ số Bình thường Cảnh báo
Consumer lag < 100 > 1000 (thêm worker)
Message/giây vào Khớp tốc độ scraper đẩy tác vụ Tăng đột biến = nghẽn ở phía sau
Message/giây ra Khớp tốc độ vào Giảm dần so với vào = worker là nút thắt

Câu hỏi thường gặp

Dưới 1.000 tác vụ/giờ có cần Kafka không?

Không hẳn. Kafka đáng dùng khi cần replay tin nhắn, thông lượng trên 100K message/giây, hoặc scale consumer group. Khối lượng nhỏ hơn thì Redis hoặc RabbitMQ đủ dùng.

Worker Kafka gọi CaptchaAI thì giải được những loại CAPTCHA nào?

Bất kỳ loại nào CaptchaAI hỗ trợ qua in.php/res.php: reCAPTCHA v2/v3, Turnstile, Cloudflare Challenge, GeeTest v3, ảnh/OCR, grid image, BLS. hCaptcha và FunCaptcha chưa hỗ trợ — đừng route vào cùng worker.

Cần bao nhiêu thread CaptchaAI để pipeline không bị nghẽn?

Số thread giới hạn số tác vụ giải song song, không phải số CAPTCHA đã giải trong tháng. Chạy 6 worker thì cần ít nhất 6 thread: BASIC ($15/tháng, 5 thread) đủ thử nghiệm, ADVANCE ($90/tháng, 50 thread) phù hợp khi scale lên vài chục worker.

Worker crash giữa chừng thì có giải trùng một CAPTCHA không?

Có thể, nếu commit offset trước khi publish kết quả. Chỉ consumer.commit() sau khi captcha-results nhận message thành công, rồi lọc trùng bằng task_id ở consumer kết quả.

CAPTCHA không giải được (poison message) thì xử lý sao?

Đặt giới hạn retry trong worker; vượt giới hạn thì publish sang topic captcha-dead-letter để kiểm tra thủ công, tránh chặn cả partition.

Xử lý sự cố thường gặp

Vấn đề Nguyên nhân Cách xử lý
Consumer lag tăng liên tục Worker không theo kịp tốc độ tác vụ Thêm worker (tối đa bằng số partition)
Kết quả bị trùng Worker crash trước khi commit offset Kiểm tra idempotency trên task_id ở consumer kết quả
Rebalance quá thường xuyên Worker crash hoặc restart liên tục Tăng session.timeout.ms; kiểm tra out-of-memory
Tác vụ phân bổ không đều Key phân phối kém Dùng key ngẫu nhiên hơn hoặc tăng số partition

Bài viết liên quan

  • Truyền phát hàng loạt kết quả CAPTCHA

Bước tiếp theo

Lấy API key CaptchaAI rồi nối Kafka vào pipeline theo đúng ba bước ở trên.

Hướng dẫn liên quan:

Os comentários estão desativados para este artigo.