DevOps và Mở Rộng

Apache Kafka + CaptchaAI: Xử lý tác vụ CAPTCHA trực tuyến

Khi khối lượng giải CAPTCHA đạt tới hàng nghìn tác vụ mỗi giờ, bạn cần nhiều hơn một hàng đợi đơn giản. Apache Kafka cung cấp tính năng truyền phát tin nhắn có độ bền cao, có trật tự, thông lượng cao — lý tưởng để tách việc gửi tác vụ CAPTCHA khỏi quá trình xử lý kết quả trên quy mô lớn.

Kiến trúc

[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 chủ đề Kafka có mối quan tâm riêng biệt:

  • captcha-tasks – các tham số CAPTCHA đang chờ giải quyết
  • captcha-results – Mã thông báo đã giải quyết sẵn sàng để sử dụng tiếp theo

Điều kiện tiên quyết

# Python
pip install kafka-python requests

# Node.js
npm install kafkajs axios

Nhà môi giới Kafka chạy trên localhost:9092 (hoặc địa chỉ cụm của bạn).

Bước 1: Tạo chủ đề

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 phân vùng cho phép tối đa sáu người tiêu dùng song song trên mỗi nhóm.

Bước 2: Trình tạo tác vụ (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();
})();

Bước 3: Công nhân CAPTCHA (Người tiêu dùng + Người giải)

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

Công nhân nhân rộng

Nhóm người tiêu dùng Kafka tự động phân phối phân vùng giữa các công nhân:

# 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

Mở rộng quy mô lên đến số lượng phân vùng. Ngoài ra, thêm nhiều phân vùng.

Giám sát

Theo dõi các số liệu chính thông qua độ trễ của người tiêu dùng Kafka:

kafka-consumer-groups.sh --describe --group captcha-workers \
  --bootstrap-server localhost:9092
Số liệu khỏe mạnh Cảnh báo
Độ trễ của người tiêu dùng < 100 > 1000 (thêm công nhân)
Tin nhắn/sec trong Phù hợp với tỷ lệ cạp Gai cho thấy sự bùng nổ
Tin nhắn/sec out Trận đấu về tỷ lệ Bị tụt lại phía sau = nút thắt

Quy tắc xác nhận của nhà sản xuất

  • Từ chối các thư thiếu loại bộ giải, siêu dữ liệu đích hoặc chi tiết định tuyến phản hồi trước khi chúng đề cập đến chủ đề.
  • Thêm khóa bình thường để các lần thử lại không tạo ra các lần giải trùng lặp cho cùng một thử thách.
  • Gửi các bản ghi không đúng định dạng đến một đường dẫn không có chữ cái có đủ ngữ cảnh để phát lại sau này.

Khắc phục sự cố

Vấn đề Nguyên nhân Cách xử lý
Độ trễ của người tiêu dùng ngày càng tăng Công nhân không thể theo kịp tốc độ nhiệm vụ Thêm nhiều phiên bản công nhân hơn (tối đa số lượng phân vùng)
Kết quả trùng lặp Công nhân gặp sự cố trước khi thực hiện bù đắp Thêm kiểm tra idempotency trên task_id trong kết quả của người tiêu dùng
Tái cân bằng quá thường xuyên Công nhân bị tai nạn/restarting Tăng session.timeout.ms; kiểm tra OOM
Nhiệm vụ không được phân bổ đồng đều Phân phối khóa kém Sử dụng khóa ngẫu nhiên hoặc nhiều phân vùng

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

Tại sao Kafka thay vì Redis hay RabbitMQ?

Kafka lý tưởng khi bạn cần độ bền của tin nhắn (khả năng phát lại), thông lượng cao (100K+ tin nhắn/sec) và mở rộng quy mô nhóm người tiêu dùng. Đối với các thiết lập đơn giản hơn dưới 1.000 tác vụ/hour, Redis hoặc RabbitMQ là đủ.

Tôi nên sử dụng một hay hai chủ đề?

Hai chủ đề (nhiệm vụ + kết quả) tách biệt rõ ràng giữa người sản xuất và người tiêu dùng. Người tạo nhiệm vụ không cần biết về người tiêu dùng kết quả và ngược lại.

Làm cách nào để xử lý các tin nhắn độc hại (CAPTCHA không thể giải quyết được)?

Đặt giới hạn thử lại trong tệp worker. Sau khi thử lại tối đa, hãy xuất bản lên chủ đề captcha-dead-letter để kiểm tra thủ công. Đừng chặn phân vùng với số lần thử lại vô hạn.

bài viết liên quan

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

Các bước tiếp theo

Xây dựng quy trình phát trực tiếp CAPTCHA —lấy khóa API CaptchaAI của bạnvà kết nối Kafka để xử lý thông lượng cao.

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

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