NATS 是一个轻量级、高性能的消息系统——没有 JVM,默认情况下没有磁盘持久性,亚毫秒级延迟。对于验证码任务分发,您需要速度和简单性而不是 Kafka 的持久性,NATS 是理想的选择。
为什么使用 NATS 执行验证码任务
| 特征 | NATS | 卡夫卡 | RabbitMQ |
|---|---|---|---|
| 延迟 | < 1 毫秒 | 5-10毫秒 | 1-5毫秒 |
| 设置复杂性 | 单一二进制 | 集群+ZooKeeper | 缓和 |
| 内存占用 | 〜20MB | 〜1GB+ | 约200MB |
| 坚持 | 可选(喷射流) | 内置 | 内置 |
| 最适合 | 临时任务,低延迟 | 持久的流媒体 | 复杂的路由 |
CAPTCHA 任务是短暂的 - 如果任务丢失,您可以重新提交。NATS 的简单性和速度使其成为自然的选择。
建筑学
[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 队列组自动在工作线程之间分发消息——每个任务都只分配给一个工作线程。
先决条件
# 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
任务发布者(爬虫)
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(队列组订阅者)
队列组确保每条消息都准确地发送给一个工作人员,即使有多个工作人员正在运行。
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);
结果收集器
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())
NATS JetStream 经久耐用
对于不能丢失的任务,启用JetStream持久化:
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}")
JetStream 增加了磁盘持久性、重放功能和一次性交付 - 与 Kafka 类似,但具有 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
故障排除
| 问题 | 原因 | 处理方式 |
|---|---|---|
| 消息被丢弃 | NATS 核心 pub/sub 不会为慢速消费者提供缓冲 | 使用 JetStream 实现持久性,或增加消费者容量 |
| 工作人员未收到消息 | 队列组名称或主题错误 | 验证主题和队列组匹配发布者 |
| 连接重置 | NATS服务器重启 | 在客户端选项中启用自动重新连接 |
| 分布不均匀 | 一名工人处理速度比其他人快 | 正常 - NATS 分配给可用的工作人员;速度更快的工作人员获得更多 |
常问问题
什么时候应该使用 NATS 而不是 Redis 或 Kafka?
当您需要最少的基础设施(单个二进制文件,无依赖项)、亚毫秒级延迟并且不需要持久消息存储时,请使用 NATS。使用 Kafka 进行持久流处理,使用 Redis 进行缓存 + pub/sub 组合。
NATS 每小时可以处理 10,000 多个验证码任务吗?
容易地。 NATS 每秒处理数百万条消息。瓶颈将是 CaptchaAI 求解时间,而不是 NATS 吞吐量。
我需要 JetStream 吗?
只要你需要坚持。对于大多数 CAPTCHA 工作流程,核心 NATS 就足够了——如果任务丢失,您可以重新提交。启用 JetStream 来实现审计跟踪或一次性处理要求。
下一步
使用 NATS 分发 CAPTCHA 任务 –获取您的 CaptchaAI API 密钥并启动轻量级工人。
相关指南:
- 卡夫卡流媒体集成
- 用于分布式处理的 Redis 队列
- RabbitMQ 消息队列集成