Chato é melhor: 67 mil eventos de telemetria por segundo no Postgres
Nesta página
Coletamos uma quantidade absurda de dados de telemetria. Prompts, métricas de uso, dados de sessão, insights de cada assistente de código com IA e de cada app de desktop que os engenheiros dos nossos clientes usam. Tudo isso escorrendo de milhares de usuários pra nossa plataforma de insights.
Na Flowstate, usamos OpenTelemetry pra medir gasto com IA. Cada chamada de API, cada sessão de código, cada invocação de modelo, ligada de volta a times, projetos e centros de custo. O volume é grande e nunca para de chegar.
Isso significa que precisamos guardar muito dado. Não é dado de “dashboard de analytics”. Não é dado de “relatório mensal”. É cada span, cada métrica, cada linha de log de cada usuário de cada cliente, gravados pra nossa pipeline de análise mastigar depois.
Todo mundo tem o mesmo primeiro instinto: “Você precisa de um data warehouse.”
Eu tentei. Tentei de verdade.
Dando uma olhada no mercado
Pra deixar claro, os produtos que avaliei são todos excelentes naquilo pra que foram feitos. São feitos pra organizações com times grandes de dados e cargas analíticas complexas. Nosso problema era mais simples: gravar muito dado de telemetria rápido, consultar depois. Pra isso, a maioria dessas soluções era mais do que a gente precisava.
Tinha também uma preocupação prática que não saía da minha cabeça. Eu deveria conseguir desenvolver dentro de um avião. Não que eu faça isso de fato, mas é a essência do problema. Se um componente da minha stack exige internet pra funcionar, não consigo construir, testar e iterar localmente. Não consigo subir no Docker e jogar dado nele num sábado de manhã. Muitos desses data warehouses pareciam resolver problemas que um Postgres bem ajustado dava conta.
O BigQuery foi a primeira parada. Ótimo pra consultar em escala, menos ideal pra escrita contínua em alto volume. A API de streaming insert cobra por linha, e isso soma rápido no nosso volume, e o perfil de latência não serve pra ingestão quase em tempo real. Um motor de análise fenomenal, só não era o encaixe certo pra uma pipeline de telemetria com muita escrita.
O Snowflake é uma plataforma de dados poderosa, mas o modelo de preço é complexo (créditos de computação, armazenamento, transferência de dados) e, pro que a gente precisava, era muito mais infraestrutura do que o problema pedia. Quando a alternativa é um Postgres bem ajustado, a comparação de custo é brutal.
O Databricks é impressionante se você tem uma área dedicada de engenharia de dados. A arquitetura lakehouse e a integração com Spark são muito boas de verdade. Mas a gente não tem um time de doze engenheiros de dados, e trazer uma plataforma dessa complexidade pro que é basicamente “grava linha, consulta linha” parecia levar um drone militar pra uma briga de faca.
Amazon Redshift, Azure Synapse e outras ofertas analíticas gerenciadas são produtos sólidos, mas cada uma traz uma carga operacional que não combinava com o ponto em que a gente estava. Mais um sistema pra monitorar, mais um conjunto de credenciais, mais um fornecedor pra lidar.
O ClickHouse foi a opção mais tentadora. Orientado a colunas, feito sob medida pra cargas analíticas com muita escrita, open source e rápido de verdade. Gostei bastante. O motivo de eu ter largado não foi a tecnologia, foi o deploy. Colocar o ClickHouse pra rodar de forma confiável num ambiente gerenciado deu mais atrito do que eu queria. O ClickHouse Cloud existe, mas aí é mais uma dependência de fornecedor, e eu já tenho Postgres gerenciado no Google Cloud SQL. Além disso, dá pra fazer consulta analítica no Postgres mesmo, então o ClickHouse parecia só mais etapas.
O que me levou de volta ao banco que eu já estava rodando.
Já construí umas coisas malucas em cima do Postgres ao longo dos anos, e a pergunta sempre aparece: “Será que um banco só basta?” Posso dizer por experiência que um banco só aguenta muita pancada. Tem gente que roda Postgres em escala de petabyte. Dá trabalho, mas dá. Então vamos ver o que conseguimos fazer no nosso caso humilde.
Eu já tenho Postgres. Ele já funciona. Vamos descobrir onde ele para de funcionar.
As falhas
O que vem a seguir é uma linha do tempo resumida de eu aprendendo as coisas do jeito difícil.
Tentativa 1: gravar direto no banco de produção
A primeira versão era tão burra quanto parece. Toda vez que chegava um evento de telemetria, a gente disparava um INSERT no Postgres de produção. Uma linha por vez. Individualmente. Que nem mandar carta pelo correio.
Funcionou em volume baixo. Funcionou em volume médio. Parou de funcionar às 3h da manhã de uma terça-feira, quando o serviço de proxy pegou um pico de tráfego e de repente a gente estava fazendo 10.000 INSERTs individuais por segundo num banco que também tentava atender a aplicação de verdade.
Eu literalmente não conseguia pegar um lock de escrita pra rodar uma migration. O pool de conexões estava saturado, cada conexão disponível estava travada num INSERT, e o runner de migration ficou ali esperando um lock que nunca vinha. Acabei matando conexões na mão pra migration passar. Às 3h da manhã. Numa terça-feira.
Tentativa 2: batching (esquentando)
A correção óbvia: parar de gravar uma linha por vez. Guardar os eventos em memória e descarregar no banco em lotes de 500 a 1.000 linhas, com INSERTs de múltiplas linhas.
Melhorou dramaticamente. Em vez de 10.000 transações, você tem de 10 a 20 transações maiores. O banco respira. O pool de conexões esvazia. As migrations rodam.
Mas apareceu um problema novo: o que acontece se a aplicação cair entre um lote e outro? Você perde o que estava no buffer.
Tentativa 3: Redis como buffer
Então colocamos o Redis como buffer intermediário. Os eventos chegam, são anexados a um Redis Stream (XADD), e um worker em segundo plano consome o stream (XREADGROUP) e descarrega os lotes no Postgres.
Isso funcionou de verdade, e o Redis continua na arquitetura final. Mas não como armazenamento de dados. Ele guarda ponteiros pros objetos que chegam, controla o que já chegou e o que já foi descarregado, e deixa a gente fazer batching de forma inteligente sem a aplicação precisar guardar estado. Se a mangueira de incêndio tem uma pressão tão alta que arrancaria a sua pele, o Redis é a válvula que transforma isso em algo manejável.
A sacada foi não usar o Redis como banco de dados. Ele é uma camada de coordenação. Os dados de telemetria em si vão direto pro Postgres. O Redis só diz o que está esperando e o que já foi processado.
A pergunta
Nesse ponto eu tinha um sistema que funcionava. Postgres pra armazenamento, Redis pra coordenação, escritas em lote com COPY pra vazão. Mas eu não sabia quanto o Postgres aguentava. Será que ele segurava 10x a carga atual? 100x? Qual é o teto de verdade?
Só tem um jeito de descobrir.
Vamos só… testar
Decidi fazer o que qualquer pessoa razoável faria: subir o Postgres num container Docker no meu notebook e jogar quantidades cada vez mais absurdas de telemetria até alguma coisa quebrar.
O setup:
- Postgres 16 rodando no Docker com 512MB de shared buffers e 500 conexões máximas
- pgBouncer em modo de transaction pooling pros testes de pool de conexões
- Um schema no estilo OpenTelemetry: tabelas
spans,metricselogscom atributos em JSONB - Um gerador de telemetria que cria traces realistas: cerca de 1.900 eventos por usuário simulado, com cerca de 3KB cada, espalhados por cerca de 100 traces com 5 a 12 spans, métricas e logs
- Quatro estratégias de escrita, testadas com números crescentes de usuários
- Nos testes extremos (10K+ usuários), rodamos rajadas de 30 segundos, porque com cerca de 200MB/s de escrita sustentada o meu SSD de 512GB encheria rápido. Mais que isso e eu estaria fazendo benchmark do modo de falha do armazenamento do meu notebook, não do Postgres.
Tudo isso rodando num MacBook Air M2 de 2022 com 24GB de RAM e SSD de 512GB. Não é um servidor. É um notebook.
As quatro estratégias
INSERT ingênuo: um INSERT por evento de telemetria. 50 conexões simultâneas. A linha de base do “por favor, não faça isso”, e exatamente o que eu fazia em produção às 3h da manhã naquela terça.
INSERT em lote: acumula eventos em lotes de 500 linhas e grava num único INSERT de múltiplas linhas. 20 conexões num pool.
Protocolo COPY: o Postgres tem um protocolo nativo de carga em massa chamado COPY. Ele joga dados separados por tabulação direto na tabela, sem passar pelo parser de SQL. É o que as ferramentas de ETL usam. Lotes de 5.000 linhas passando pelo pg-copy-streams.
Pool + lote: a mesma estratégia de lotes, mas roteada pelo pgBouncer em modo de transaction pooling. Testa se o pool de conexões acrescenta vazão de verdade com muita concorrência.
Os números
Com 1.000 usuários simulados, cada um gerando cerca de 1.900 eventos de cerca de 3KB, a gente grava aproximadamente 1,8 milhão de linhas de telemetria realista em três tabelas. Payloads JSONB de verdade, com nomes de modelo, contagem de tokens, dados de custo, IDs de sessão. O tipo de dado que você veria numa pipeline de telemetria de IA em produção.
Veja o que aconteceu.
Escritas por segundo
Em volume baixo, tudo parece mais ou menos igual. Você não percebe a diferença entre as estratégias quando está gravando um par de milhares de linhas.
Mas aumente a escala e as linhas se separam. A abordagem ingênua estaciona em cerca de 14.000 escritas/s e fica por lá. O gargalo é o custo por instrução e a disputa por conexões.
O protocolo COPY chega a 67.000 escritas por segundo. Num notebook. Num container Docker. Com max_wal_size em 4GB e as configurações padrão de checkpoint. Com payloads de 3KB, isso dá cerca de 200MB/s de vazão de escrita sustentada. Se a gente tivesse ajustado os timeouts de checkpoint e a compressão do WAL, provavelmente dava pra ir mais longe.
Uma ressalva sobre o COPY: ou tudo ou nada. Se uma linha de um lote de 5.000 tiver JSON malformado ou violar uma constraint, o lote inteiro falha. O INSERT em lote deixa você usar cláusulas ON CONFLICT pra lidar com duplicatas e dados sujos sem drama. Em produção, a gente valida antes do COPY e cai pro INSERT em lote quando algo parece suspeito.
O INSERT em lote fica no meio, com cerca de 34K escritas/s. Sólido e prático. Você não precisa aprender uma API nova pra usar.
A estratégia com pool pelo pgBouncer ficou entre 43 e 49K escritas/s, mais rápida que o lote puro porque o pgBouncer reaproveita conexões melhor que o pool da nossa aplicação.
Frente a frente com 1.000 usuários
Com 1.000 usuários (1,8M de eventos), o COPY faz 67K escritas/s enquanto o ingênuo fica preso em 14K. É quase 5x mais rápido pros mesmos dados. Nesses tamanhos de payload, o COPY gravou 916K eventos em 13,6 segundos. O ingênuo levou mais de um minuto.
Mas o que me surpreendeu foi isto: nenhuma das estratégias caiu. Eu esperava que o Postgres começasse a engasgar. Erros de conexão, OOM kills, WAL inchado. Nada disso aconteceu. O Postgres simplesmente aguentou.
Então, naturalmente, decidi subir o volume.
Um choque de realidade
Vamos ser claros: hoje a gente não tem 10 milhões de usuários simultâneos mandando 1.900 eventos de telemetria por dia. Se tivesse, estaria ganhando dinheiro suficiente pra contratar um departamento inteiro pra se preocupar com arquitetura de banco de dados.
Hoje, o nosso volume real de produção é tranquilamente absorvido pelo nosso Postgres atual, sem suar. Então por que empurrar a simulação até milhões de usuários gerando centenas de megabytes por segundo?
Estritamente pra provar um ponto.
Eu queria saber o que acontece quando a escolha “chata” finalmente bate na parede. E o que descobri é que, mesmo em escalas absurdas e teóricas em que o Postgres começa a ficar lento e as consultas pioram, ele não simplesmente morre. Ele degrada com elegância. E, mais importante, quando ele fica lento, existem alavancas padrão e chatas que você pode puxar pra recuperar a velocidade.
O teste de verdade: escritas E leituras ao mesmo tempo
Benchmark de escrita isolado mente pra você.
Em produção, você não pode pausar as leituras enquanto ingere dados. Nossa pipeline de análise roda consultas agregadas nesses dados enquanto eles estão sendo gravados. Cálculos de percentil em milhões de spans. Mapas de dependência entre serviços. Quebras de taxa de erro. Consultas que fazem o Postgres pensar de verdade.
Então fiz o óbvio: rodei escritas COPY em streaming e 8 consultas analíticas ao mesmo tempo, escalando de 1.000 usuários até 10 milhões. Nos testes maiores, usei janelas de rajada de 30 segundos, porque com essas taxas de escrita o SSD de 512GB do meu MacBook encheria fisicamente em menos de 10 minutos.
Com 1.000 usuários (escrita completa, 1,8M de eventos), o COPY sustentou 39K escritas/s enquanto atendia, ao mesmo tempo, mais de 7.000 consultas analíticas com p95 de 9ms. A consulta mais lenta, o JOIN do mapa de dependências, levou 1,2 segundo.
Com 10 milhões de usuários simulados, numa rajada de 30 segundos, o Postgres manteve 16.763 escritas/s enquanto atendia consultas analíticas com p95 abaixo de 200ms. Em 30 segundos, ele gravou 514.536 linhas de telemetria de 3KB. Em todos os pontos de escala, a consulta do mapa de dependências (um self-join na tabela de spans) foi consistentemente a mais lenta, chegando a cerca de 1,2 segundo.
A vazão de escrita cai conforme a tabela cresce e as leituras disputam I/O. Isso é esperado. Mas o Postgres nunca caiu, nunca estourou a memória, nunca corrompeu nada. Ficou mais lento e seguiu em frente. Toda consulta devolveu o resultado correto. Toda escrita foi confirmada.
E lembre: isso rodou com só 512MB de shared_buffers. O banco era bem maior que o cache. Se eu tivesse dado 4GB de shared_buffers neste Mac de 24GB, o conjunto de trabalho inteiro caberia na RAM e essas consultas seriam bem mais rápidas. Deliberadamente não fiz isso, porque banco de produção nem sempre consegue manter tudo na memória.
Tá, mas dá pra consertar as leituras lentas?
Aquela consulta do mapa de dependências estava me incomodando. Então adicionei índices e rodei o teste de novo em cerca de 920K linhas (a telemetria de 500 usuários).
Sete índices pontuais: busca por nome de serviço, filtros de status code, joins por trace ID, relações de span pai, buscas compostas de métricas e distribuição de severidade. Depois um ANALYZE pra atualizar as estatísticas do planejador de consultas.
O resultado de destaque:
Consulta de erros recentes: de 51ms pra 1,2ms. Um ganho de 43x. Pra deixar claro, status_code e start_time são colunas de primeiro nível no schema, não ficam enterradas no payload JSONB. A coluna JSONB attributes guarda as coisas flexíveis (headers HTTP, contagem de tokens, nomes de modelo, dados de custo). Os campos que a gente consulta com frequência são colunas extraídas, com tipos de verdade. O índice parcial em status_code = 2 combinado com o índice em start_time DESC deixou o Postgres pular o full table scan por completo. Ele só caminha pelo índice e pega os 50 primeiros.
Distribuição de severidade de log: de 65ms pra 33ms. O índice composto em (severity, service_name) transforma um sequential scan num index-only scan. 2x mais rápido.
A consulta do mapa de dependências caiu de 456ms pra 320ms. Uma melhora de 1,4x. Melhor, mas nada transformador. Essa consulta faz um self-join em centenas de milhares de linhas, e nenhum índice sozinho elimina o custo fundamental de correlacionar spans pai e filho nessa escala.
Mas tudo bem. Existem caminhos claros de escala pra quando você precisar.
O roteiro de escala (pra quando você realmente precisar)
Esta é a parte em que vou defender que você não otimize cedo demais. O que a gente construiu funciona. Aguenta a nossa carga atual com folga, e aguenta 10x sem suar. Mas se um dia chegarmos a processar centenas de milhões de eventos, há dois caminhos claros, ambos ainda usando Postgres.
Caminho 1: réplicas de leitura
O movimento de escala mais simples. Monte uma réplica com streaming replication e aponte todas as consultas analíticas pra ela. As escritas vão pro primário, as leituras vão pra réplica. O primário se concentra inteiramente na ingestão, e a réplica pode mastigar JOINs complexos sem afetar a vazão de escrita.
É uma mudança de terça à tarde. Google Cloud SQL, AWS RDS e Azure suportam réplicas de leitura nativamente. Você adiciona uma connection string e uma regra de roteamento. Sua vazão de escrita volta pra faixa de 67K escritas/s do banco isolado, porque ele não está mais brigando com as leituras, e suas consultas analíticas podem demorar o quanto quiserem na réplica sem ninguém perceber. Em hardware de servidor de verdade, com mais núcleos de CPU e armazenamento mais rápido, esse número seria bem maior.
No nosso caso, em que a análise não precisa ser em tempo real, uma réplica com alguns segundos de atraso de replicação é perfeitamente aceitável.
Caminho 2: particionamento de tabelas
O particionamento por tempo divide suas tabelas em pedaços. Uma partição por dia, por semana ou por mês. Consultas que filtram por tempo só varrem as partições relevantes, em vez da tabela inteira. Aquela consulta de 1,2 segundo do mapa de dependências? Se você olha só as últimas 24 horas em vez do histórico completo, está varrendo uma fração dos dados. A consulta cai de segundos pra milissegundos.
O particionamento também deixa o ciclo de vida dos dados trivial. Quer apagar dados com mais de 90 dias? DROP TABLE spans_2025_q4. Sem vacuum, sem bloat, sem tabela travada. Instantâneo.
A opção de escala gigante
Se um dia a gente precisasse ir pra um nível realmente maluco, com centenas de milhões de usuários e bilhões de eventos de telemetria, o Postgres também tem resposta pra isso. O Citus é uma extensão do Postgres (totalmente open source, da Microsoft) que distribui seus dados por vários nós de Postgres usando sharding por hash. Você roda CREATE EXTENSION citus; e pronto. Seu schema continua o mesmo. Suas consultas continuam as mesmas. Você só tem mais nós fazendo o trabalho.
Faça o shard da tabela de spans por trace_id e cada nó cuida de uma fatia do conjunto total. Uma consulta por um trace específico cai em um nó. Uma consulta agregada se espalha por todos os nós e junta os resultados. Escala horizontal linear sem sair do ecossistema Postgres. Mesmas ferramentas, mesmo monitoramento, mesma experiência do time.
A gente não precisa disso. Provavelmente não vai precisar por muito tempo. Mas o fato de o caminho existir, sem sair do Postgres, é exatamente o motivo de começar simples ter sido a decisão certa.
Não otimize cedo demais
Quero ser bem claro sobre isto: a arquitetura que a gente roda hoje não é a arquitetura mais otimizada pra esse problema. É a mais apropriada.
Temos uma pipeline de dados. Conseguimos rodar consultas complexas enquanto despejamos escritas nela. Ela aguenta cargas analíticas simultâneas. Não cai. Então por que eu precisaria de outro produto de banco de dados pra isso?
Claro, o I/O de disco pode virar um limite em algum momento. Mas existem opções de escala bem definidas, e sabemos exatamente quais são, porque medimos onde ficam os gargalos. Essa é a vantagem de começar simples: quando precisar escalar, você sabe o que escalar.
A gente poderia ter começado com ClickHouse. Poderia ter montado o Citus desde o primeiro dia. Poderia ter construído uma arquitetura Lambda com Kafka, um processador de stream, uma camada de serving, uma camada de batch e uma… você entendeu.
Mas toda essa complexidade tem custo. Cada sistema a mais é mais uma coisa pra monitorar, mais uma coisa que pode quebrar às 3h da manhã, mais uma coisa que o seu time precisa entender. Se você não tem escala que justifique, está pagando o imposto da complexidade sem receber o benefício da escalabilidade.
Postgres no Cloud SQL, com Redis como camada de coordenação, usando o protocolo COPY pra ingestão em lote. Essa é a nossa arquitetura. Ela aguenta a nossa carga atual. Vai aguentar 10x a nossa carga atual. E quando ela não aguentar 100x, a gente vai saber exatamente onde estão os gargalos, porque já mediu.
A boa notícia é que a gente não precisa que esse serviço seja rápido. Os dados de telemetria são processados em lotes pra análise, não servidos em tempo real pros usuários. Uma consulta que leva 4 segundos em vez de 40 milissegundos é totalmente aceitável no nosso caso.
Mas se um dia a gente precisasse que fosse rápido… bom, você viu os números. Tem muita folga.
O que eu realmente aprendi
- Nunca use INSERTs individuais pra escrita de alta vazão. Lote ou COPY. Sempre.
- O protocolo COPY existe por um motivo. Não é só pra carga inicial de dados. É uma estratégia legítima de ingestão em produção.
- Pool de conexões não é opcional em escala. pgBouncer em modo transaction. Sem desculpa.
- Teste escritas E leituras juntas. Benchmark só de escrita engana. O perfil de desempenho real é o que acontece quando o seu banco está fazendo as duas coisas ao mesmo tempo.
- Índices importam, mas não pra tudo. Um índice parcial bem posicionado pode dar um ganho de 43x em consultas específicas. Mas consultas agregadas em milhões de linhas vão ser lentas de qualquer jeito. É aí que você apela pra réplicas e particionamento.
- Não otimize cedo demais. Comece com a arquitetura mais simples que funciona. Meça onde estão os gargalos. Escale as partes que precisam escalar. Nem tudo, nem tudo de uma vez.
- O Postgres aguenta mais do que as pessoas imaginam. 67K escritas/s de eventos de telemetria de 3KB. 39K escritas/s enquanto atende consultas analíticas ao mesmo tempo. Escalado pra 10M de usuários simulados e ainda assim não caiu. Num notebook. No Docker.
- Teste suas suposições. Passei semanas lendo posts de blog comparando data warehouses. Eu poderia ter passado uma tarde com Docker e um script e ter números de verdade. Os números contaram uma história melhor.
Ah, e mais uma coisa. A gente passou este post inteiro colocando o Postgres num inferno absoluto. Milhões de linhas, leituras concorrentes, self-joins em spans gigantescos. E nem pensamos no Redis na frente dele. O Redis que coordena tudo isso, que lida com a mangueira de incêndio da telemetria chegando, que controla ponteiros e gerencia o estado dos lotes. Ele nem piscou. A gente não fez benchmark dele porque não havia o que medir. Ele simplesmente funcionou.
A maioria dos problemas dá pra resolver com um servidor web, um Redis e um Postgres.
O código do benchmark é escrito em Go e está em github.com/willhackett/bench-postgres. docker-compose up -d, go run . -mode=full pras estratégias de escrita, go run . -mode=chaos pro teste de caos com escrita e leitura combinadas e go run . -mode=optimize pra comparação de indexação.