DevOps & Scaling

NATS Messaging + CaptchaAI:轻量级验证码任务分发

验证码任务量一大,很多团队第一反应是上 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,跑起来就能开始识别任务。

相关指南:

该文章已禁用评论。