Pesquisar este blog
Automação, desenvolvimento em Python e C#, e soluções práticas para otimizar sua rotina profissional e pessoal. Documentando códigos e aprendizados reais.
Postagem em destaque
- Gerar link
- X
- Outros aplicativos
Ferramentas MCP: Da coleta de dados à LLM Wiki
Coleta de dados em ambientes heterogêneos: a dor de quem precisa de um pipeline confiável
O desenvolvedor que precisa alimentar um LLM com informações de múltiplas fontes (APIs internas, bancos legados, arquivos CSV na nuvem e logs de máquinas IoT) costuma perder horas apenas garantindo que o scrape não quebre ao mudar um endpoint ou ao encontrar um registro corrompido. O resultado costuma ser um pipeline que funciona em teste, mas explode em produção por falta de idempotência, tratamento de falhas e monitoramento.
Arquitetura modular baseada em MCP (Message, Collect, Process)
Dividir o fluxo em três estágios claros evita acoplamento e permite escalar cada parte independentemente:
- Message: camada de transporte (Kafka, RabbitMQ, ou Google Pub/Sub) que garante entrega ao menos uma vez e permite replay.
- Collect: workers responsáveis por buscar dados de origem. Cada worker deve ser stateless e configurável via JSON/YAML.
- Process: normalização, enriquecimento e gravação em um repositório de documentos (Elasticsearch, Pinecone ou um bucket S3 pronto para ingestão no LLM).
Escolha do broker
Kafka oferece retenção configurável e compactação de tópicos, ideal para replay de falhas de Collect. Em ambientes menores, RabbitMQ com plugins de dead‑letter é suficiente e tem menor overhead operacional.
Worker de coleta genérico em Python
import json
import httpx
from typing import Any, Dict, Iterable
from pathlib import Path
class Collector:
def __init__(self, config_path: str) -> None:
self.cfg = json.loads(Path(config_path).read_text())
self.client = httpx.AsyncClient(timeout=30)
async def fetch(self, endpoint: str, params: Dict[str, Any]) -> Dict[str, Any]:
try:
resp = await self.client.get(endpoint, params=params)
resp.raise_for_status()
return resp.json()
except httpx.HTTPError as exc:
# Log estruturado para o broker de mensagens
raise RuntimeError(f"Fetch failed for {endpoint}: {exc}") from exc
async def run(self) -> Iterable[Dict[str, Any]]:
for src in self.cfg["sources"]:
data = await self.fetch(src["url"], src.get("params", {}))
# Normaliza campos críticos
yield {
"source_id": src["id"],
"payload": data,
"fetched_at": src["timestamp"]
}
# Uso típico dentro de um container Docker
# collector = Collector("/etc/collector/config.json")
# async for record in collector.run():
# await producer.send("raw-data", record)
Observações:
- Uso de
httpx.AsyncClientgarante alta concorrência sem bloquear a thread. - Erros são propagados como
RuntimeErrorpara que o broker possa redirecionar a mensagem para um tópico de dead‑letter. - O formato de saída já inclui metadados que facilitam o rastreamento de origem.
Processamento idempotente e enriquecimento
O estágio Process deve ser capaz de detectar duplicatas e aplicar transformações determinísticas. A estratégia mais segura é usar um hash do payload como chave de deduplicação.
import hashlib
import json
from typing import Dict, Any
def payload_hash(payload: Dict[str, Any]) -> str:
# Serializa com ordenação de chaves para garantir hash estável
canonical = json.dumps(payload, sort_keys=True, separators=(",", ":"))
return hashlib.sha256(canonical.encode()).hexdigest()
def enrich(record: Dict[str, Any]) -> Dict[str, Any]:
# Exemplo de enriquecimento: lookup de taxonomia em um cache Redis
from redis import Redis
r = Redis(host="redis", port=6379, db=0)
tax_id = record["payload"].get("category_id")
tax_name = r.get(f"tax:{tax_id}") or "unknown"
record["payload"]["category_name"] = tax_name.decode()
record["hash"] = payload_hash(record["payload"])
return record
Ao gravar no destino, use o hash como document_id. Sistemas como Elasticsearch ou Pinecone rejeitam documentos com ID já existente, garantindo idempotência automática.
Persistência para LLM Wiki: escolha entre vetores e texto bruto
Um LLM que responde perguntas sobre a "Wiki interna" precisa de duas camadas:
- Texto completo indexado para busca lexical (BM25, Elasticsearch).
- Embeddings vetoriais para recuperação semântica (FAISS, Milvus, ou Pinecone).
Manter os dois índices sincronizados evita que a resposta contenha trechos desatualizados.
Pipeline de ingestão final
import asyncio
from elasticsearch import AsyncElasticsearch
from sentence_transformers import SentenceTransformer
es = AsyncElasticsearch(hosts=["http://es:9200"])
model = SentenceTransformer("all-MiniLM-L6-v2")
async def index_document(doc: Dict[str, Any]) -> None:
# 1. Indexa texto bruto
await es.index(
index="wiki-raw",
id=doc["hash"],
document={"content": doc["payload"]["text"], "metadata": doc["payload"]}
)
# 2. Gera embedding e salva em vetor store (exemplo com Pinecone)
vector = model.encode(doc["payload"]["text"]).tolist()
# pinecone.upsert([(doc["hash"], vector, {"source": doc["source_id"]})])
# O código acima é ilustrativo; substitua por client específico.
# Orquestração simples
async def main():
async for raw in producer.consume("processed-data"):
await index_document(raw)
asyncio.run(main())
Notas de produção:
- Separar índices permite usar políticas de retenção diferentes (texto por 2 anos, vetores por 6 meses).
- O modelo de embeddings deve ser carregado uma única vez por worker; reutilizar a mesma instância reduz latência em 30‑40%.
- Monitorar
es.indexe o cliente de vetor para erros de timeout; reprocessar mensagens falhas via tópico de retry.
Recomendações práticas e armadilhas frequentes
- Não misture lógica de coleta com persistência. Cada camada deve ter seu próprio container e seu próprio ciclo de vida.
- Use schemas versionados. Quando um campo novo aparece, adicione‑o ao schema sem remover o antigo; isso evita que mensagens antigas quebrem o
Process. - Dead‑letter e replay são obrigatórios. Configure tópicos de retry com back‑off exponencial e limite de tentativas para evitar loops infinitos.
- Teste de carga realista. Simule picos de 10‑20x a taxa média de ingestão; verifique latência de
httpx, throughput do broker e tempo de indexação no Elasticsearch. - Evite dependências pesadas no worker. Bibliotecas de embeddings podem consumir GB de RAM; isole-as em containers dedicados ou use serviços gerenciados.
- Log estruturado e correlação de IDs. Inclua
source_id,hashetrace_idem todos os logs para rastrear falhas de ponta a ponta.
Checklist de implantação
- Broker configurado com retenção mínima de 48h e tópicos de dead‑letter.
- Containers de
Collectorcom limites de memória < 512Mi e auto‑scale baseado em CPU. - Instância Elasticsearch com shards adequados ao volume esperado (pelo menos 1 shard por 10 GB de texto).
- Modelo de embeddings carregado em serviço separado ou em GPU se a taxa de ingestão exigir.
- Alertas de latência > 2 s em
fetch,indexeembeddingvia Prometheus/Grafana.
- Gerar link
- X
- Outros aplicativos
Postagens mais visitadas
RPA em Escala: Criando uma Fila de Trabalho para seus Robôs com Python e Redis Queue (RQ)
- Gerar link
- X
- Outros aplicativos
Como preparar dados para aplicações com LLMs
- Gerar link
- X
- Outros aplicativos
Comentários
Postar um comentário