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:
- O scraper publica a tarefa no subject
captcha.tasks. - O queue group
captcha-workersentrega essa mensagem a exatamente um worker disponível. - O worker envia a tarefa para a API da CaptchaAI e aguarda a resolução.
- O worker publica o resultado no subject
captcha.results. - Um coletor consome
captcha.resultse 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 parain.php/res.phpe 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: