验证码任务量一大,很多团队第一反应是上 Kafka——结果 ZooKeeper、集群运维全跟着来了。任务能重试时根本不需要这么重:NATS 一个二进制文件即可跑,亚毫秒延迟。下面搭一套 NATS 验证码分发系统:发布、识别、收结果,顺带看要不要上 JetStream。
为什么验证码任务适合用 NATS,而不是 Kafka / RabbitMQ
验证码任务的三个特点,决定了它更适合轻量消息系统:
- 体积小——几百字节的 JSON
- 生命短——几十秒内就该出结果
- 丢了直接重新提交,无需“绝对不丢”
NATS 比 Kafka 的持久化更划算:
| 特性 | NATS | Kafka | RabbitMQ |
|---|---|---|---|
| 延迟 | < 1 毫秒 | 5-10 毫秒 | 1-5 毫秒 |
| 部署复杂度 | 单个二进制文件 | 集群 + ZooKeeper | 中等 |
| 内存占用 | ~20 MB | ~1 GB+ | ~200 MB |
| 持久化 | 可选(JetStream) | 内置 | 内置 |
| 最适合 | 临时任务、低延迟场景 | 持久化流处理 | 复杂路由 |
NATS 的简单和速度,正好匹配“一次性、丢了就重发”这种场景。
整体架构:爬虫发布任务,worker 调用 CaptchaAI 识别
[Scrapers] → Publish → [NATS: captcha.tasks]
↓
Queue Group: captcha-workers
├── Worker 1 (solve via CaptchaAI)
├── Worker 2
└── Worker 3
↓
Publish → [NATS: captcha.results]
↓
[Result Subscribers]
NATS 的队列组(queue group)自动把消息分给多个 worker,不用自己写负载均衡。
典型场景:
- 多站点、多验证码类型混跑(reCAPTCHA v2、Cloudflare Turnstile)
- 用一个
captcha.tasks主题汇总 - worker 按
method字段路由,无需单独搭分发系统
环境准备
# Install NATS server
# macOS
brew install nats-server
# Linux
curl -L https://github.com/nats-io/nats-server/releases/download/v2.10.0/nats-server-v2.10.0-linux-amd64.tar.gz | tar xz
# Start
nats-server
# Python client
pip install nats-py
# Node.js client
npm install nats
国内网络慢时:pip 加
-i https://pypi.tuna.tsinghua.edu.cn/simple走清华镜像,NATS server 直接docker run -p 4222:4222 nats拉镜像启动。
第一步:任务发布端(爬虫写入队列)
把验证码参数打包成 JSON 推到 captcha.tasks 主题:
Python
import asyncio
import json
import nats
async def publish_captcha_tasks():
nc = await nats.connect("nats://localhost:4222")
tasks = [
{
"task_id": f"task_{i}",
"method": "userrecaptcha",
"sitekey": "6Le-wvkSAAAAAPBMRTvw0Q4Muexq9bi0DJwx_mJ-",
"pageurl": f"https://example.com/page/{i}"
}
for i in range(100)
]
for task in tasks:
await nc.publish("captcha.tasks", json.dumps(task).encode())
print(f"Published: {task['task_id']}")
await nc.flush()
await nc.close()
asyncio.run(publish_captcha_tasks())
JavaScript
const { connect, StringCodec } = require("nats");
const sc = StringCodec();
async function publishCaptchaTasks() {
const nc = await connect({ servers: "nats://localhost:4222" });
for (let i = 0; i < 100; i++) {
const task = {
task_id: `task_${i}`,
method: "userrecaptcha",
sitekey: "6Le-wvkSAAAAAPBMRTvw0Q4Muexq9bi0DJwx_mJ-",
pageurl: `https://example.com/page/${i}`,
};
nc.publish("captcha.tasks", sc.encode(JSON.stringify(task)));
console.log(`Published: ${task.task_id}`);
}
await nc.flush();
await nc.close();
}
publishCaptchaTasks();
第二步:CAPTCHA Worker(队列组订阅者)
队列组能保证同一条消息只会分给一个 worker,哪怕同时跑了好几个实例也不会重复处理。
Python
import asyncio
import json
import os
import nats
import aiohttp
API_KEY = os.environ["CAPTCHAAI_API_KEY"]
async def solve_captcha(session, task):
"""Submit to CaptchaAI and poll for result."""
# Submit
async with session.post("https://ocr.captchaai.com/in.php", data={
"key": API_KEY,
"method": task["method"],
"googlekey": task["sitekey"],
"pageurl": task["pageurl"],
"json": 1
}) as resp:
data = await resp.json(content_type=None)
if data.get("status") != 1:
return {"task_id": task["task_id"], "error": data.get("request")}
captcha_id = data["request"]
# Poll for result
for _ in range(60):
await asyncio.sleep(5)
async with session.get("https://ocr.captchaai.com/res.php", params={
"key": API_KEY, "action": "get", "id": captcha_id, "json": 1
}) as resp:
result = await resp.json(content_type=None)
if result.get("status") == 1:
return {"task_id": task["task_id"], "solution": result["request"]}
if result.get("request") != "CAPCHA_NOT_READY":
return {"task_id": task["task_id"], "error": result.get("request")}
return {"task_id": task["task_id"], "error": "TIMEOUT"}
async def worker(worker_id):
nc = await nats.connect("nats://localhost:4222")
# Subscribe with queue group — each message goes to one worker only
sub = await nc.subscribe("captcha.tasks", queue="captcha-workers")
print(f"Worker {worker_id} listening...")
async with aiohttp.ClientSession() as session:
async for msg in sub.messages:
task = json.loads(msg.data.decode())
print(f"Worker {worker_id} processing {task['task_id']}")
result = await solve_captcha(session, task)
# Publish result
await nc.publish(
"captcha.results",
json.dumps(result).encode()
)
status = "solved" if "solution" in result else result.get("error")
print(f" → {task['task_id']}: {status}")
asyncio.run(worker(1))
JavaScript
const { connect, StringCodec } = require("nats");
const axios = require("axios");
const sc = StringCodec();
const API_KEY = process.env.CAPTCHAAI_API_KEY;
function sleep(ms) {
return new Promise((r) => setTimeout(r, 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 { task_id: task.task_id, 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 { task_id: task.task_id, solution: result.data.request };
}
if (result.data.request !== "CAPCHA_NOT_READY") {
return { task_id: task.task_id, error: result.data.request };
}
}
return { task_id: task.task_id, error: "TIMEOUT" };
}
async function worker(workerId) {
const nc = await connect({ servers: "nats://localhost:4222" });
// Queue group subscription — load-balanced across workers
const sub = nc.subscribe("captcha.tasks", { queue: "captcha-workers" });
console.log(`Worker ${workerId} listening...`);
for await (const msg of sub) {
const task = JSON.parse(sc.decode(msg.data));
console.log(`Worker ${workerId} processing ${task.task_id}`);
const result = await solveCaptcha(task);
nc.publish("captcha.results", sc.encode(JSON.stringify(result)));
const status = result.solution ? "solved" : result.error;
console.log(` → ${task.task_id}: ${status}`);
}
}
worker(1);
第三步:收集识别结果
收集端订阅 captcha.results,按 solved / failed 分类计数:
async def collect_results():
nc = await nats.connect("nats://localhost:4222")
sub = await nc.subscribe("captcha.results")
solved = 0
failed = 0
async for msg in sub.messages:
result = json.loads(msg.data.decode())
if "solution" in result:
solved += 1
print(f"[SOLVED] {result['task_id']} — {result['solution'][:30]}...")
else:
failed += 1
print(f"[FAILED] {result['task_id']} — {result['error']}")
print(f" Stats: {solved} solved, {failed} failed")
asyncio.run(collect_results())
量大了以后建议把 solved / failed 计数接入 Prometheus,别只靠终端打印。
需要持久化?打开 JetStream
任务不能丢时,用 JetStream 把 NATS 升级成带持久化、可重放的消息系统:
async def durable_publisher():
nc = await nats.connect("nats://localhost:4222")
js = nc.jetstream()
# Create stream (one-time setup)
await js.add_stream(name="CAPTCHA", subjects=["captcha.>"])
# Publish with acknowledgment
ack = await js.publish("captcha.tasks", json.dumps(task).encode())
print(f"Published to stream, seq={ack.seq}")
相当于给 NATS 加上 Kafka 级别的持久化,同时保留轻量和简单。
横向扩展 worker
吞吐不够就加 worker,NATS 自动摊分任务:
# Run multiple workers — NATS distributes automatically via queue groups
python worker.py --id=1 &
python worker.py --id=2 &
python worker.py --id=3 &
# Each task goes to exactly one worker
# Add more workers to increase throughput
分发规则:“谁先处理完谁先拿下一条”,不是按 worker 数量平均分——快的自然拿得多,属于正常现象。
常见故障排查
| 问题 | 原因 | 处理方式 |
|---|---|---|
| 消息丢失 | 核心 pub/sub 不为慢消费者缓冲 | 换 JetStream 做持久化,或扩容 worker |
| worker 收不到消息 | 队列组名或 subject 拼错 | 确认 subject、队列组名与发布端一致 |
| 连接被重置 | NATS server 重启 | 客户端配置里打开自动重连 |
| 各 worker 处理量不均 | 有的处理得快 | 正常——快的自然拿得多 |
常见问题
worker 掉线,正在处理的任务会丢失吗?
会丢,核心 pub/sub 不缓冲消息:
- 任务廉价,重新提交即可
- 接受不了丢任务,上 JetStream 加 ack
什么场景该选 NATS,而不是 Redis 或 Kafka?
- 要最小基础设施、亚毫秒延迟:选 NATS
- 要持久化流处理:选 Kafka
- 只想要缓存 + pub/sub 组合:选 Redis
国内网络下装 nats-py、npm nats 慢怎么办?
参考上面「环境准备」的镜像方案:pip/npm 换清华源,NATS server 直接拉 Docker 镜像。
下一步
申请 CaptchaAI API Key,跑起来就能开始识别任务。
相关指南:
- Kafka 流式处理集成
- 用 Redis 队列做分布式处理
- RabbitMQ 消息队列集成