DevOps 与扩展

Redis Queue + CaptchaAI:分布式验证码处理

Redis 提供快速、可靠的队列,用于在多个工作人员之间分配验证码任务。本指南构建了一个完整的生产者-消费者系统,具有结果跟踪和错误处理功能。


建筑学

Producers → Redis List (FIFO queue) → Workers → CaptchaAI API
                                          ↓
                                   Redis Hash (results)
                                          ↓
                                   Redis Pub/Sub (notifications)

任务队列管理器

import json
import time
import uuid
import redis


class CaptchaQueue:
    """Redis-backed CAPTCHA task queue."""

    def __init__(self, redis_url="redis://localhost:6379"):
        self.redis = redis.from_url(redis_url)
        self.queue_key = "captcha:tasks"
        self.results_key = "captcha:results"
        self.notify_channel = "captcha:done"

    def submit(self, method, params, priority="normal"):
        """Submit a CAPTCHA task to the queue."""
        task_id = str(uuid.uuid4())[:8]
        task = {
            "id": task_id,
            "method": method,
            "params": params,
            "submitted_at": time.time(),
            "priority": priority,
        }

        if priority == "high":
            self.redis.lpush(self.queue_key, json.dumps(task))
        else:
            self.redis.rpush(self.queue_key, json.dumps(task))

        return task_id

    def fetch(self, timeout=30):
        """Fetch next task from queue (blocking)."""
        result = self.redis.blpop(self.queue_key, timeout=timeout)
        if result is None:
            return None
        _, raw = result
        return json.loads(raw)

    def store_result(self, task_id, result):
        """Store task result and notify listeners."""
        self.redis.hset(
            self.results_key,
            task_id,
            json.dumps(result),
        )
        # Notify via pub/sub
        self.redis.publish(self.notify_channel, task_id)
        # Set TTL on result (1 hour)
        # Results are in a hash, so we track expiry separately
        self.redis.setex(
            f"captcha:ttl:{task_id}", 3600, "1",
        )

    def get_result(self, task_id):
        """Get result for a task (non-blocking)."""
        raw = self.redis.hget(self.results_key, task_id)
        if raw:
            return json.loads(raw)
        return None

    def wait_result(self, task_id, timeout=120):
        """Wait for a task result via polling."""
        start = time.time()
        while time.time() - start < timeout:
            result = self.get_result(task_id)
            if result:
                return result
            time.sleep(1)
        return None

    def queue_stats(self):
        """Get queue statistics."""
        return {
            "pending": self.redis.llen(self.queue_key),
            "completed": self.redis.hlen(self.results_key),
        }

工人进程

import os
import time
import requests


class QueueWorker:
    """Worker that processes CAPTCHA tasks from Redis queue."""

    def __init__(self, api_key, queue):
        self.api_key = api_key
        self.queue = queue
        self.base = "https://ocr.captchaai.com"

    def run(self):
        """Main worker loop."""
        worker_id = os.getpid()
        print(f"Worker {worker_id} started")

        while True:
            task = self.queue.fetch(timeout=30)
            if task is None:
                continue

            task_id = task["id"]
            print(f"[{worker_id}] Processing {task_id}")

            start = time.time()
            try:
                token = self._solve(task["method"], task["params"])
                duration = time.time() - start
                self.queue.store_result(task_id, {
                    "status": "success",
                    "token": token,
                    "duration": f"{duration:.1f}s",
                })
                print(f"[{worker_id}] {task_id} solved in {duration:.1f}s")

            except Exception as e:
                self.queue.store_result(task_id, {
                    "status": "error",
                    "error": str(e),
                })
                print(f"[{worker_id}] {task_id} failed: {e}")

    def _solve(self, method, params, timeout=120):
        resp = requests.post(f"{self.base}/in.php", data={
            "key": self.api_key,
            "method": method,
            "json": 1,
            **params,
        }, timeout=30)
        result = resp.json()

        if result.get("status") != 1:
            raise RuntimeError(result.get("request"))

        captcha_id = result["request"]
        start = time.time()

        while time.time() - start < timeout:
            time.sleep(5)
            resp = requests.get(f"{self.base}/res.php", params={
                "key": self.api_key,
                "action": "get",
                "id": captcha_id,
                "json": 1,
            }, timeout=15)
            data = resp.json()
            if data["request"] != "CAPCHA_NOT_READY":
                if data.get("status") == 1:
                    return data["request"]
                raise RuntimeError(data["request"])

        raise TimeoutError("Solve timeout")


# Run worker
if __name__ == "__main__":
    queue = CaptchaQueue()
    worker = QueueWorker(os.environ["CAPTCHAAI_KEY"], queue)
    worker.run()

多工作人员启动器

import multiprocessing
import os


def start_workers(num_workers=4):
    """Launch multiple worker processes."""
    queue = CaptchaQueue()
    processes = []

    for i in range(num_workers):
        p = multiprocessing.Process(
            target=run_worker,
            args=(os.environ["CAPTCHAAI_KEY"],),
        )
        p.start()
        processes.append(p)
        print(f"Started worker {i + 1}/{num_workers}")

    return processes


def run_worker(api_key):
    queue = CaptchaQueue()
    worker = QueueWorker(api_key, queue)
    worker.run()


# Launch
processes = start_workers(num_workers=4)

生产者示例

queue = CaptchaQueue()

# Submit tasks
urls = [
    "https://site1.com/login",
    "https://site2.com/register",
    "https://site3.com/checkout",
]

task_ids = []
for url in urls:
    tid = queue.submit("userrecaptcha", {
        "googlekey": "SITE_KEY",
        "pageurl": url,
    })
    task_ids.append(tid)
    print(f"Submitted {tid} for {url}")

# Wait for all results
for tid in task_ids:
    result = queue.wait_result(tid, timeout=120)
    status = result["status"] if result else "timeout"
    print(f"{tid}: {status}")

# Check queue stats
print(queue.queue_stats())

Pub/Sub 结果监听器

任务完成后立即收到通知:

import threading


def listen_results(queue):
    """Listen for completed task notifications."""
    pubsub = queue.redis.pubsub()
    pubsub.subscribe(queue.notify_channel)

    for message in pubsub.listen():
        if message["type"] == "message":
            task_id = message["data"].decode()
            result = queue.get_result(task_id)
            print(f"Task {task_id} completed: {result['status']}")


# Run listener in background
listener = threading.Thread(
    target=listen_results,
    args=(CaptchaQueue(),),
    daemon=True,
)
listener.start()

故障排除

问题 原因 处理方式
工作人员因队列中的任务而空闲 连接到错误的 Redis 验证 REDIS_URL
结果消失 无 TTL 管理 使用 setex 使结果过期
队列无限增长 工人太慢了 添加更多工作人员或增加并发性
重复处理 任务弹出但工作线程崩溃 使用 brpoplpush 实现可靠队列

常问问题

为什么选择 Redis 而不是其他消息队列?

Redis 简单、快速,大多数团队已经在运行它。对于复杂的路由或有保证的交付,请考虑使用 RabbitMQ。

每个 Redis 实例有多少个工作线程?

单个 Redis 实例可处理 100k+ 操作/second. 瓶颈是 CaptchaAI API 吞吐量,而不是 Redis。根据您的 CaptchaAI 容量规划工作人员。

我应该使用 Redis 流而不是列表吗?

Redis Streams 提供消费者组和确认,这对于生产来说更好。列表非常适合简单的设置。


相关指南


分配你的工作量——获取CaptchaAI今天。

该文章已禁用评论。