DevOps & Scaling

Mensagens NATS + CaptchaAI: distribuição leve de tarefas CAPTCHA

Se seus scrapers publicam milhares de tarefas CAPTCHA por hora, o gargalo raramente é a fila de mensagens — é o tempo de resolução do CAPTCHA em si. Ainda assim, a escolha do broker importa: com o NATS você distribui essas tarefas entre workers com um único binário, sem cluster, sem ZooKeeper, com latência abaixo de 1 milissegundo. A diferença para o Kafka é que tarefas CAPTCHA são efêmeras — se uma se perder, o custo de reenviá-la é praticamente zero, então a durabilidade de um log distribuído raramente compensa o overhead operacional.

Por que usar NATS para tarefas CAPTCHA

Em números: o NATS entrega mensagens em menos de 1 ms, contra 5-10 ms do Kafka e 1-5 ms do RabbitMQ. A instalação é um binário único (~20 MB de memória) — sem o cluster e o ZooKeeper que o Kafka exige, e sem a configuração moderada do RabbitMQ (~200 MB). Persistência é opcional via JetStream; Kafka e RabbitMQ já nascem com persistência embutida, mas você paga esse custo mesmo quando não precisa dele. Para tarefas efêmeras e de baixa latência — publicar um CAPTCHA e esperar o worker resolver — essa simplicidade é o que conta: o que se perder, você reenvia, e segue em frente.

Arquitetura: do scraper ao worker e de volta

O fluxo completo, passo a passo:

  1. O scraper publica a tarefa no subject captcha.tasks.
  2. O queue group captcha-workers entrega essa mensagem a exatamente um worker disponível.
  3. O worker envia a tarefa para a API da CaptchaAI e aguarda a resolução.
  4. O worker publica o resultado no subject captcha.results.
  5. Um coletor consome captcha.results e soma as estatísticas de sucesso e falha.
[Scrapers] → Publish → [NATS: captcha.tasks]
                              ↓
                    Queue Group: captcha-workers
                    ├── Worker 1 (solve via CaptchaAI)
                    ├── Worker 2
                    └── Worker 3
                              ↓
                    Publish → [NATS: captcha.results]
                              ↓
                    [Result Subscribers]

O diagrama acima é a mesma lógica em forma visual: o queue group faz com que cada tarefa chegue a um único worker, mesmo com vários workers conectados ao mesmo tempo — ninguém processa a mesma tarefa duas vezes, e a carga se distribui sozinha conforme os workers ficam disponíveis.

Publicando tarefas CAPTCHA (o scraper)

Comece instalando o servidor e o client do seu stack; depois é só publicar as tarefas — o exemplo abaixo cobre Python e, na sequência, o equivalente em Node.js:

# 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
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())
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();

Workers CAPTCHA: consumindo a fila com queue groups

Cada worker se conecta ao mesmo queue group captcha-workers. O NATS cuida da distribuição: mesmo com vários workers rodando ao mesmo tempo, cada mensagem chega a um único worker — nunca a mais de um. Os exemplos abaixo enviam a tarefa para a CaptchaAI, fazem o polling do resultado e publicam a resposta em captcha.results; primeiro em Python, depois em Node.js:

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))
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);

Coletando resultados e persistindo com JetStream

Um assinante simples do subject captcha.results já é suficiente para acompanhar o que foi resolvido e o que falhou. Quando a tarefa não pode se perder — por exemplo, para trilha de auditoria — ative a persistência do JetStream, que adiciona disco, replay e entrega exatamente uma vez, mantendo a simplicidade do NATS:

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())
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}")

Na prática, a maioria dos pipelines de CAPTCHA não precisa disso: como a tarefa é barata de reenviar, o NATS Core (sem JetStream) já resolve o dia a dia.

Escalando workers e operando em produção

Escalar não exige coordenação manual — basta subir mais processos worker apontando para o mesmo queue group:

# 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

O NATS distribui automaticamente entre os workers disponíveis; o teto real de throughput não é o broker, é quantas requisições simultâneas o seu plano da CaptchaAI permite. Um pool de 10 workers, cada um mantendo até 5 requisições simultâneas em aberto, casa exatamente com o plano ADVANCE (US$ 90/mês, 50 threads); escalando para 20 workers nessas condições, o plano CORPORATE (US$ 240/mês, 150 threads) ainda tem folga. Antes de ir para produção, vale conhecer os problemas mais comuns desse tipo de pipeline:

Workers em sa-east-1 (São Paulo)? O publish/subscribe do NATS continua abaixo de 1 ms dentro da região — o tempo que pesa no pipeline é o da chamada para in.php/res.php e o tempo de resolução do CAPTCHA, não o transporte de mensagens.

  • Mensagens descartadas — o NATS Core não guarda buffer para consumidores lentos. Ative o JetStream para persistência ou aumente a capacidade dos workers.
  • Worker não recebe mensagens — subject ou nome do queue group incorretos. Confira se publisher e worker usam exatamente os mesmos valores.
  • Conexão reiniciando — o servidor NATS reiniciou. Habilite reconexão automática nas opções do client.
  • Distribuição desigual entre workers — normal: o NATS entrega para quem está disponível, e workers mais rápidos acabam recebendo mais tarefas.

Perguntas frequentes

O NATS funciona bem com workers hospedados no Brasil (AWS sa-east-1)?

Sim. Se o broker e os workers rodam na mesma região — por exemplo, sa-east-1 — o publish/subscribe continua abaixo de 1 ms. O tempo que realmente pesa no pipeline vem da chamada à API da CaptchaAI e do tempo de resolução do CAPTCHA, não do transporte de mensagens.

Quantos workers preciso para processar 10 mil tarefas CAPTCHA por hora?

Menos do que se imagina, e o dimensionamento depende do plano de threads da CaptchaAI, não do NATS — ele entrega milhões de mensagens por segundo sem esforço. Calcule pelo número de requisições simultâneas que seu plano permite: um plano STANDARD (US$ 30/mês, 15 threads) já sustenta um volume relevante se o tempo médio de resolução for baixo; para volumes maiores, suba para ADVANCE ou PREMIUM conforme a concorrência necessária.

O NATS Core resolve ou preciso do JetStream em produção?

Na maioria dos pipelines de CAPTCHA, o NATS Core já é suficiente — como as tarefas são efêmeras, perder uma e reenviá-la custa pouco. Ative o JetStream só quando precisar de trilha de auditoria, replay de mensagens ou entrega exatamente uma vez.

O que acontece se um worker cair no meio de uma tarefa?

Sem JetStream, a tarefa em processamento se perde e o scraper (ou uma rotina de retry) precisa reenviá-la — como o custo de reenvio é baixo, isso raramente é um problema real. Com JetStream habilitado, a mensagem permanece no stream até receber confirmação (ack), então um worker que reinicia retoma o trabalho em vez de descartar a tarefa.

Próximos passos

Distribua suas tarefas CAPTCHA com NATS: obtenha sua chave de API da CaptchaAI e coloque workers leves no ar em minutos. Guias relacionados:

Os comentários estão desativados para este artigo.