验证码识别任务一多,同步阻塞调用马上顶不住:一个 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_nack(requeue=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 间隔,加上断线重连逻辑 |
相关指南
- Redis 队列 + CaptchaAI 分布式处理方案
- 批量验证码识别:多任务处理指南
队列稳,识别才稳——从 CaptchaAI 开始接入 RabbitMQ。