DevOps & Scaling

RabbitMQ + CaptchaAI: Integração de fila de mensagens

Um worker que trava no meio da resolução de um CAPTCHA não pode deixar a tarefa sumir — e é exatamente esse o risco de uma fila simples. Com filas duráveis, confirmação de mensagem (ack/nack) e uma dead letter exchange para tarefas que falham repetidamente, o RabbitMQ garante que nada se perde, mesmo com reinício do broker ou queda de um worker. Veja como montar essa integração, do zero até produção.

O que você vai montar nesta página:

  • Um produtor que publica tarefas de CAPTCHA com prioridade e persistência em disco
  • Um worker consumidor que resolve o CAPTCHA pela API e confirma (ou rejeita) a mensagem
  • Um coletor de resultados que junta as respostas prontas pelo task_id
  • Roteamento por tipo de CAPTCHA, para filas e workers especializados

Por que o RabbitMQ é a fila de mensagens certa para CAPTCHA

Antes de instalar qualquer coisa, compare o que o RabbitMQ resolve que uma lista simples em Redis não resolve:

  • Filas duráveis — as tarefas sobrevivem a reinícios do broker
  • Confirmação de mensagem (ack/nack) — nenhuma tarefa se perde numa falha de worker
  • Dead letter exchange — tarefas com falha são roteadas para investigação, não descartadas
  • Filas com prioridade — CAPTCHAs urgentes furam a fila de tarefas comuns
  • Routing keys — encaminhamento por tipo de CAPTCHA para workers especializados

Se algum desses pontos já te mordeu em produção, é isso que os próximos passos resolvem.

Preparando o ambiente

Suba o broker e instale o cliente Python:

# Docker
docker run -d --hostname rabbitmq \
  -p 5672:5672 -p 15672:15672 \
  rabbitmq:3-management

# Python client
pip install pika requests

Produtor: enviando tarefas de CAPTCHA para a fila

O produtor declara as filas — incluindo a dead letter exchange — e publica cada tarefa com delivery_mode=2 para persistência em disco. A priority garante que CAPTCHAs urgentes furem a fila de tarefas comuns.

import json
import uuid
import pika


class CaptchaProducer:
    """Submit CAPTCHA tasks to RabbitMQ."""

    def __init__(self, rabbitmq_url="amqp://guest:guest@localhost:5672/"):
        self.connection = pika.BlockingConnection(
            pika.URLParameters(rabbitmq_url),
        )
        self.channel = self.connection.channel()
        self._setup_queues()

    def _setup_queues(self):
        """Declare durable queues and exchanges."""
        # Dead letter exchange for failed tasks
        self.channel.exchange_declare(
            exchange="captcha.dlx",
            exchange_type="direct",
            durable=True,
        )
        self.channel.queue_declare(
            queue="captcha.failed",
            durable=True,
        )
        self.channel.queue_bind(
            queue="captcha.failed",
            exchange="captcha.dlx",
            routing_key="failed",
        )

        # Main task queue with dead letter routing
        self.channel.queue_declare(
            queue="captcha.tasks",
            durable=True,
            arguments={
                "x-dead-letter-exchange": "captcha.dlx",
                "x-dead-letter-routing-key": "failed",
                "x-message-ttl": 300000,  # 5 min TTL
            },
        )

        # Results queue
        self.channel.queue_declare(
            queue="captcha.results",
            durable=True,
        )

    def submit(self, method, params, priority=0):
        """Submit a CAPTCHA task."""
        task_id = str(uuid.uuid4())[:8]
        task = {
            "id": task_id,
            "method": method,
            "params": params,
        }

        self.channel.basic_publish(
            exchange="",
            routing_key="captcha.tasks",
            body=json.dumps(task),
            properties=pika.BasicProperties(
                delivery_mode=2,  # Persistent
                priority=priority,
                message_id=task_id,
            ),
        )
        return task_id

    def close(self):
        self.connection.close()


# Usage
producer = CaptchaProducer()

task_id = producer.submit("userrecaptcha", {
    "googlekey": "SITE_KEY",
    "pageurl": "https://example.com",
}, priority=5)

print(f"Submitted: {task_id}")
producer.close()

O submit() devolve um task_id curto, usado para rastrear a tarefa até o coletor de resultados, mais adiante.

Consumidor: o worker que resolve o CAPTCHA

Cada worker processa uma tarefa por vez (prefetch_count=1), para não travar numa resolução longa enquanto outras tarefas esperam. Ele envia o CAPTCHA para in.php e consulta res.php a cada 5 segundos até receber o token ou estourar o timeout.

import json
import os
import time
import pika
import requests


class CaptchaConsumer:
    """RabbitMQ consumer that solves CAPTCHAs."""

    def __init__(self, api_key, rabbitmq_url="amqp://guest:guest@localhost:5672/"):
        self.api_key = api_key
        self.base = "https://ocr.captchaai.com"
        self.connection = pika.BlockingConnection(
            pika.URLParameters(rabbitmq_url),
        )
        self.channel = self.connection.channel()
        # Process one task at a time
        self.channel.basic_qos(prefetch_count=1)

    def start(self):
        """Start consuming tasks."""
        self.channel.basic_consume(
            queue="captcha.tasks",
            on_message_callback=self._handle_task,
        )
        print("Worker started. Waiting for tasks...")
        self.channel.start_consuming()

    def _handle_task(self, ch, method, properties, body):
        """Process a single CAPTCHA task."""
        task = json.loads(body)
        task_id = task["id"]
        print(f"Processing {task_id}...")

        try:
            token = self._solve(task["method"], task["params"])

            # Publish result
            result = {
                "task_id": task_id,
                "status": "success",
                "token": token,
            }
            ch.basic_publish(
                exchange="",
                routing_key="captcha.results",
                body=json.dumps(result),
                properties=pika.BasicProperties(delivery_mode=2),
            )

            # Acknowledge message (remove from queue)
            ch.basic_ack(delivery_tag=method.delivery_tag)
            print(f"{task_id} solved successfully")

        except Exception as e:
            print(f"{task_id} failed: {e}")

            # Reject and send to dead letter queue
            ch.basic_nack(
                delivery_tag=method.delivery_tag,
                requeue=False,  # Goes to DLX
            )

    def _solve(self, captcha_method, params, timeout=120):
        resp = requests.post(f"{self.base}/in.php", data={
            "key": self.api_key,
            "method": captcha_method,
            "json": 1,
            **params,
        }, timeout=30)
        result = resp.json()

        if result.get("status") != 1:
            raise RuntimeError(result.get("request"))

        captcha_id = result["request"]
        start = time.time()

        while time.time() - start < timeout:
            time.sleep(5)
            resp = requests.get(f"{self.base}/res.php", params={
                "key": self.api_key,
                "action": "get",
                "id": captcha_id,
                "json": 1,
            }, timeout=15)
            data = resp.json()
            if data["request"] != "CAPCHA_NOT_READY":
                if data.get("status") == 1:
                    return data["request"]
                raise RuntimeError(data["request"])

        raise TimeoutError("Solve timeout")


# Run worker
if __name__ == "__main__":
    consumer = CaptchaConsumer(os.environ["CAPTCHAAI_KEY"])
    consumer.start()

Se a chamada falhar, basic_nack com requeue=False manda a tarefa direto para a dead letter exchange, sem retentativa automática nesse ponto — o que evita que um CAPTCHA malformado consuma threads do plano em loop. Vale reparar na região dos seus workers: se eles rodam longe do broker — uma instância em AWS sa-east-1 (São Paulo), por exemplo, enquanto o restante do pipeline está em outra região —, a latência de rede até https://ocr.captchaai.com entra no tempo medido pelo timeout, e comparar tempos de resolução entre regiões sem isolar esse RTT costuma levar a conclusões erradas.

Coletor de resultados

Depois que o worker publica o resultado na fila captcha.results, um processo separado consome essas mensagens e as associa de volta ao pedido original pelo task_id.

import json
import pika


class ResultCollector:
    """Collect task results from the results queue."""

    def __init__(self, rabbitmq_url="amqp://guest:guest@localhost:5672/"):
        self.connection = pika.BlockingConnection(
            pika.URLParameters(rabbitmq_url),
        )
        self.channel = self.connection.channel()
        self.results = {}

    def collect(self, expected_count, timeout=120):
        """Collect a specific number of results."""
        deadline = time.time() + timeout

        while len(self.results) < expected_count and time.time() < deadline:
            method, _, body = self.channel.basic_get(
                queue="captcha.results",
                auto_ack=True,
            )
            if body:
                result = json.loads(body)
                self.results[result["task_id"]] = result

            time.sleep(0.5)

        return self.results

Roteamento por tipo: filas especializadas por worker

Nem todo CAPTCHA tem o mesmo custo de resolução. Separar reCAPTCHA, Turnstile e imagem em filas próprias permite escalar cada worker de acordo com o tipo que ele processa, em vez de misturar tudo numa fila genérica:

# Setup exchanges and queues
channel.exchange_declare(
    exchange="captcha.types",
    exchange_type="direct",
    durable=True,
)

# Queue per type
for captcha_type in ["recaptcha", "turnstile", "image"]:
    channel.queue_declare(queue=f"captcha.{captcha_type}", durable=True)
    channel.queue_bind(
        queue=f"captcha.{captcha_type}",
        exchange="captcha.types",
        routing_key=captcha_type,
    )


# Submit with routing
def submit_routed(channel, captcha_type, task):
    channel.basic_publish(
        exchange="captcha.types",
        routing_key=captcha_type,
        body=json.dumps(task),
        properties=pika.BasicProperties(delivery_mode=2),
    )

Checklist antes de ir para produção

  1. durable=True em todas as filas e delivery_mode=2 em toda publicação — sem isso, uma queda do broker apaga a fila inteira.
  2. prefetch_count=1 por worker, para uma resolução de CAPTCHA lenta não travar as tarefas seguintes.
  3. Dead letter exchange configurada e monitorada — declarar não basta se ninguém olha captcha.failed.
  4. Intervalo de heartbeat definido, com lógica de reconexão no cliente para sobreviver a quedas de rede.
  5. Chave de API da CaptchaAI em variável de ambiente ou gerenciador de segredos, nunca embutida no código do worker.

Erros comuns na fila de mensagens RabbitMQ

Problema Causa Correção
Mensagens perdidas numa queda Fila não durável durable=True e delivery_mode=2
Worker travado numa única tarefa Resolução de CAPTCHA demorada prefetch_count=1 por worker
Dead letter queue crescendo Falhas persistentes Revise as tarefas com falha e corrija os parâmetros
Conexão caindo Timeout de heartbeat Configure o intervalo de heartbeat e adicione lógica de reconexão

Ao investigar uma dead letter queue crescendo, evite logar o corpo completo das mensagens em produção: sitekeys e tokens não deveriam parar em log agregado por padrão. Trate esse dado como sensível sob a LGPD (ou o RGPD, em Portugal).

Perguntas frequentes

Quando vale trocar o Redis pelo RabbitMQ nas filas de CAPTCHA?

Quando você precisa de confirmação de mensagem, dead letter routing ou roteamento por tipo de CAPTCHA. Para filas simples com latência mínima, o Redis continua sendo a opção mais leve.

Quantos workers devo rodar por núcleo de CPU?

Um consumidor por núcleo costuma bastar. Como cada worker processa uma tarefa por vez (prefetch_count=1), 4 núcleos equivalem a 4 workers em paralelo.

Como descubro por que uma tarefa caiu na dead letter queue?

Consuma a fila captcha.failed e inspecione o payload original. Na maioria dos casos o motivo é sitekey incorreta, pageurl inacessível ou saldo insuficiente na conta — o campo request retornado pela CaptchaAI costuma indicar o erro exato.

Faz sentido usar filas com prioridade para CAPTCHAs urgentes?

Sim, quando parte do tráfego é sensível a tempo — um checkout aguardando o token em tempo real, por exemplo — e o restante pode esperar. Use o parâmetro priority no basic_publish para esses casos, mas evite marcar tudo como prioritário: isso anula o benefício.

É seguro guardar a chave de API da CaptchaAI direto no código do worker?

Não. Use uma variável de ambiente (CAPTCHAAI_KEY, como no exemplo desta página) ou um gerenciador de segredos. Isso evita que a chave vaze em um repositório ou em logs de erro do worker.

Guias relacionados

Filas confiáveis para CAPTCHA — comece com a CaptchaAI e o RabbitMQ.

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