DevOps & Scaling

RabbitMQ + CaptchaAI:消息队列集成

验证码识别任务一多,同步阻塞调用马上顶不住:一个 reCAPTCHA 挑战卡 20 秒,10 个并发请求就能拖死采集脚本。把提交、识别、结果分发拆开,用消息队列串起来——这正是 RabbitMQ 要解决的问题。

本文用 RabbitMQ + CaptchaAI 搭一套生产级识别管道:持久化队列保证任务不丢,死信交换自动收集失败任务,优先队列让紧急验证码插队,按类型路由分给不同 worker。代码可直接复制运行。

读完本文你会拿到:

  • 持久化生产者/worker 代码
  • 按类型路由到独立队列的完整示例

验证码识别为什么要接 RabbitMQ 消息队列

单进程轮询在量小时没问题,并发一上来,任务丢失、worker 卡死这些坑迟早会踩到:

能力 解决的问题
持久化队列 Broker 重启后,未处理的任务不会丢
消息确认(ack) worker 中途崩溃时,任务会重新投递,不会静默丢失
死信交换(DLX) 失败任务自动转入专门队列,方便集中排查
优先级队列 紧急验证码可以插队,不用排在批量任务后面
路由键(routing key) 按验证码类型把任务分给专门的 worker

环境准备:启动 RabbitMQ 并安装依赖

# Docker
docker run -d --hostname rabbitmq \
  -p 5672:5672 -p 15672:15672 \
  rabbitmq:3-management

# Python client
pip install pika requests

国内网络下装依赖慢怎么办

国内网络装 PyPI 包偶尔会慢,可加镜像参数:pip install -i https://pypi.tuna.tsinghua.edu.cn/simple pika requests


生产者:把验证码任务发进队列

import json
import uuid
import pika


class CaptchaProducer:
    """Submit CAPTCHA tasks to RabbitMQ."""

    def __init__(self, rabbitmq_url="amqp://guest:guest@localhost:5672/"):
        self.connection = pika.BlockingConnection(
            pika.URLParameters(rabbitmq_url),
        )
        self.channel = self.connection.channel()
        self._setup_queues()

    def _setup_queues(self):
        """Declare durable queues and exchanges."""
        # Dead letter exchange for failed tasks
        self.channel.exchange_declare(
            exchange="captcha.dlx",
            exchange_type="direct",
            durable=True,
        )
        self.channel.queue_declare(
            queue="captcha.failed",
            durable=True,
        )
        self.channel.queue_bind(
            queue="captcha.failed",
            exchange="captcha.dlx",
            routing_key="failed",
        )

        # Main task queue with dead letter routing
        self.channel.queue_declare(
            queue="captcha.tasks",
            durable=True,
            arguments={
                "x-dead-letter-exchange": "captcha.dlx",
                "x-dead-letter-routing-key": "failed",
                "x-message-ttl": 300000,  # 5 min TTL
            },
        )

        # Results queue
        self.channel.queue_declare(
            queue="captcha.results",
            durable=True,
        )

    def submit(self, method, params, priority=0):
        """Submit a CAPTCHA task."""
        task_id = str(uuid.uuid4())[:8]
        task = {
            "id": task_id,
            "method": method,
            "params": params,
        }

        self.channel.basic_publish(
            exchange="",
            routing_key="captcha.tasks",
            body=json.dumps(task),
            properties=pika.BasicProperties(
                delivery_mode=2,  # Persistent
                priority=priority,
                message_id=task_id,
            ),
        )
        return task_id

    def close(self):
        self.connection.close()


# Usage
producer = CaptchaProducer()

task_id = producer.submit("userrecaptcha", {
    "googlekey": "SITE_KEY",
    "pageurl": "https://example.com",
}, priority=5)

print(f"Submitted: {task_id}")
producer.close()

priority 参数如何影响处理顺序

submit() 只负责把任务塞进队列并返回 task_id,不等待识别结果。priority 越高,越先被取走,可用来把紧急任务和批量任务分开。


消费者 worker:从队列取任务并调用 CaptchaAI 识别

import json
import os
import time
import pika
import requests


class CaptchaConsumer:
    """RabbitMQ consumer that solves CAPTCHAs."""

    def __init__(self, api_key, rabbitmq_url="amqp://guest:guest@localhost:5672/"):
        self.api_key = api_key
        self.base = "https://ocr.captchaai.com"
        self.connection = pika.BlockingConnection(
            pika.URLParameters(rabbitmq_url),
        )
        self.channel = self.connection.channel()
        # Process one task at a time
        self.channel.basic_qos(prefetch_count=1)

    def start(self):
        """Start consuming tasks."""
        self.channel.basic_consume(
            queue="captcha.tasks",
            on_message_callback=self._handle_task,
        )
        print("Worker started. Waiting for tasks...")
        self.channel.start_consuming()

    def _handle_task(self, ch, method, properties, body):
        """Process a single CAPTCHA task."""
        task = json.loads(body)
        task_id = task["id"]
        print(f"Processing {task_id}...")

        try:
            token = self._solve(task["method"], task["params"])

            # Publish result
            result = {
                "task_id": task_id,
                "status": "success",
                "token": token,
            }
            ch.basic_publish(
                exchange="",
                routing_key="captcha.results",
                body=json.dumps(result),
                properties=pika.BasicProperties(delivery_mode=2),
            )

            # Acknowledge message (remove from queue)
            ch.basic_ack(delivery_tag=method.delivery_tag)
            print(f"{task_id} solved successfully")

        except Exception as e:
            print(f"{task_id} failed: {e}")

            # Reject and send to dead letter queue
            ch.basic_nack(
                delivery_tag=method.delivery_tag,
                requeue=False,  # Goes to DLX
            )

    def _solve(self, captcha_method, params, timeout=120):
        resp = requests.post(f"{self.base}/in.php", data={
            "key": self.api_key,
            "method": captcha_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__":
    consumer = CaptchaConsumer(os.environ["CAPTCHAAI_KEY"])
    consumer.start()

prefetch_count 与失败重试的关系

prefetch_count=1 避免慢任务堵死整个 worker。失败时 basic_nackrequeue=False)转进死信队列 captcha.failed


结果收集:从 captcha.results 队列拿回 token

import json
import pika


class ResultCollector:
    """Collect task results from the results queue."""

    def __init__(self, rabbitmq_url="amqp://guest:guest@localhost:5672/"):
        self.connection = pika.BlockingConnection(
            pika.URLParameters(rabbitmq_url),
        )
        self.channel = self.connection.channel()
        self.results = {}

    def collect(self, expected_count, timeout=120):
        """Collect a specific number of results."""
        deadline = time.time() + timeout

        while len(self.results) < expected_count and time.time() < deadline:
            method, _, body = self.channel.basic_get(
                queue="captcha.results",
                auto_ack=True,
            )
            if body:
                result = json.loads(body)
                self.results[result["task_id"]] = result

            time.sleep(0.5)

        return self.results

轮询还是订阅:怎么选

collect() 用轮询方式从 captcha.results 拿结果,适合"提交一批、等全部回来"的批量场景;如果是长驻服务,用 basic_consume 订阅这个队列会更省资源。


按类型路由:给 reCAPTCHA、Turnstile、GeeTest 分配专属队列

国内站点常见 GeeTest(极验)滑块验证码,海外站点以 reCAPTCHA、Cloudflare Turnstile 为主,混用一个队列不好排查。按类型拆队列更实用:GeeTest v3(已支持,v4 仅"即将支持")走独立 worker;reCAPTCHA、Turnstile 依赖 Google/Cloudflare 托管脚本,国内网络访问不太稳定,单独部署更清晰。

将不同的验证码类型路由给专门的 worker:

# Setup exchanges and queues
channel.exchange_declare(
    exchange="captcha.types",
    exchange_type="direct",
    durable=True,
)

# Queue per type
for captcha_type in ["recaptcha", "turnstile", "image"]:
    channel.queue_declare(queue=f"captcha.{captcha_type}", durable=True)
    channel.queue_bind(
        queue=f"captcha.{captcha_type}",
        exchange="captcha.types",
        routing_key=captcha_type,
    )


# Submit with routing
def submit_routed(channel, captcha_type, task):
    channel.basic_publish(
        exchange="captcha.types",
        routing_key=captcha_type,
        body=json.dumps(task),
        properties=pika.BasicProperties(delivery_mode=2),
    )

拆完之后两类 worker 可以分别扩容,互不拖累。


常见问题

RabbitMQ 和 Redis 队列该怎么选?

需要保证传递、死信路由或按类型路由用 RabbitMQ;简单低延迟场景 Redis 更省事,两者都能配合 CaptchaAI。

BASIC 套餐($15/月,5 线程)能跑起这套架构吗?

能,套餐等级不影响这套架构;worker 多了再升级到 STANDARD($30/月,15 线程)。

国内环境下识别 GeeTest,要不要单独建队列?

建议单独建。GeeTest(极验)只连 CaptchaAI 接口,国内网络下比较稳定;reCAPTCHA、Turnstile 依赖 Google/Cloudflare 资源,分开更好排查。

死信队列(DLQ)堆积怎么排查,能自动重试吗?

偶发几条,查一下 sitekey/pageurl 参数;持续增长多半是参数错或余额不足。可给死信队列配 TTL 延迟重试交换,拒收的消息会自动转回主队列重试。

worker 重启时,正在处理的任务会不会被重复识别?

有小概率会。ack 之前 worker 挂掉,消息会重新投递给别的 worker。多数场景不影响结果,介意的话可按 task_id 做幂等去重。


排错清单

先按这张表定位问题:

问题 原因 处理方式
崩溃后消息丢失 队列或消息没设置持久化 队列设 durable=True,消息设 delivery_mode=2
单个 worker 卡在一个任务上 没限制单次领取数量,慢任务占着 worker 每个 worker 设 prefetch_count=1
死信队列一直变多 任务持续失败 核对 sitekey/pageurl 参数,检查账户余额
连接总掉线 心跳超时 设置 heartbeat 间隔,加上断线重连逻辑

相关指南


队列稳,识别才稳——从 CaptchaAI 开始接入 RabbitMQ。

该文章已禁用评论。