O aborrecido é melhor: 67 000 eventos de telemetria por segundo no Postgres
Nesta página
Recolhemos uma quantidade enorme de dados de telemetria. Prompts, métricas de utilização, dados de sessão, informação de todos os assistentes de código com IA e aplicações de desktop que os engenheiros dos nossos clientes usam. Tudo a afunilar de milhares de utilizadores para a nossa plataforma de insights.
Na Flowstate, usamos OpenTelemetry para medir os gastos com IA. Cada chamada à API, cada sessão de código, cada invocação de um modelo, ligada às equipas, aos projetos e aos centros de custo. O volume é grande e nunca pára de chegar.
Isto significa que precisamos de guardar muitos dados. Não dados de “dashboard de analítica”. Não dados de “relatório mensal”. Cada span, cada métrica, cada linha de log de cada utilizador de cada cliente, persistidos para o nosso pipeline de análise os poder mastigar mais tarde.
O primeiro instinto de toda a gente é o mesmo: “Precisas de um data warehouse.”
Tentei. Tentei a sério.
À procura de alternativas
Para ser claro, os produtos que analisei são todos excelentes naquilo para que foram feitos. Foram feitos para organizações com equipas de dados grandes e cargas analíticas complexas. O nosso problema era mais simples: escrever muitos dados de telemetria depressa, consultá-los depois. Para isso, a maioria destas soluções era mais do que precisávamos.
Havia também uma preocupação prática que não me largava. Devo conseguir desenvolver num avião. Não é que o faça mesmo, mas é a essência do problema. Se uma peça da minha stack precisa de internet para funcionar, não consigo compilar, testar e iterar localmente. Não consigo pô-la a correr em Docker e atirar-lhe dados num sábado de manhã. Muitos destes data warehouses pareciam tentar resolver problemas que um Postgres bem afinado resolvia.
O BigQuery foi a primeira paragem. Ótimo para consultas à escala, menos bom para escritas contínuas de grande volume. A API de streaming inserts tem custos por linha que se acumulam depressa ao nosso volume, e o perfil de latência não serve para ingestão quase em tempo real. Um motor de analítica fenomenal, mas não é o indicado para um pipeline de telemetria pesado em escritas.
O Snowflake é uma plataforma de dados poderosa, mas o modelo de preços é complexo (créditos de computação, armazenamento, transferência de dados) e, para o que precisávamos, era muito mais infraestrutura do que o problema pedia. Quando a alternativa é um Postgres bem afinado, a comparação de custos é brutal.
O Databricks impressiona se tiveres uma função dedicada de engenharia de dados. A arquitetura lakehouse e a integração com o Spark são mesmo poderosas. Mas nós não temos uma equipa de doze engenheiros de dados, e meter uma plataforma daquela complexidade para o que no fundo é “escrever linhas, ler linhas” soou-me a levar um drone militar para uma luta de facas.
O Amazon Redshift, o Azure Synapse e outras ofertas analíticas geridas são todos bons produtos, mas cada um traz uma carga operacional que não condizia com a fase em que estávamos. Mais um sistema para monitorizar, mais um conjunto de credenciais, mais uma relação com um fornecedor.
O ClickHouse foi a opção mais convincente. Orientado a colunas, feito de propósito para cargas analíticas pesadas em escritas, open source e rápido a sério. Gostei muito dele. Se o larguei, não foi por causa da tecnologia, foi por causa da instalação. Pô-lo a funcionar de forma fiável num ambiente gerido deu mais atrito do que eu queria. O ClickHouse Cloud existe, mas é mais uma dependência de um fornecedor quando já tenho um Postgres gerido no Google Cloud SQL. E de qualquer forma dá para fazer consultas analíticas em Postgres, por isso o ClickHouse parecia-me passos a mais.
O que me trouxe de volta à base de dados que eu já tinha.
Ao longo dos anos construí coisas bem malucas em Postgres, e a pergunta aparece sempre: “Temos a certeza de que uma base de dados chega?” Garanto-te por experiência que uma base de dados aguenta muita pancada. Há quem corra Postgres à escala do petabyte. Faz-se, com dificuldade, mas faz-se. Por isso vamos ver o que conseguimos fazer para o nosso caso humilde.
Já tenho Postgres. Já funciona. Vamos descobrir onde deixa de funcionar.
As falhas
O que se segue é uma cronologia condensada de eu a aprender as coisas à bruta.
Tentativa 1: escrever diretamente na base de dados de produção
A primeira versão era tão estúpida quanto parece. Sempre que chegava um evento de telemetria, disparávamos um INSERT contra a instância de produção do Postgres. Uma linha de cada vez. Individualmente. Como quem posta cartas.
Funcionou com volume baixo. Funcionou com volume médio. Deixou de funcionar às 3 da manhã de uma terça-feira, quando o serviço de proxy apanhou um pico de tráfego e de repente estávamos a fazer 10 000 INSERTs individuais por segundo contra uma base de dados que também tentava servir a aplicação a sério.
Não conseguia, literalmente, obter um bloqueio de escrita para correr uma migração. O pool de ligações estava saturado, todas as ligações disponíveis estavam presas num INSERT, e o executor de migrações ficou ali à espera de um bloqueio que nunca ia chegar. Acabei a matar ligações à mão para fazer passar a migração. Às 3 da manhã. Numa terça-feira.
Tentativa 2: batching (a aquecer)
A correção óbvia: deixar de escrever uma linha de cada vez. Guardar os eventos em memória e despejá-los na base de dados em lotes de 500 a 1 000 linhas com INSERTs de várias linhas.
Foi dramaticamente melhor. Em vez de 10 000 transações, ficas com 10 a 20 transações maiores. A base de dados respira. O pool de ligações esvazia. As migrações correm.
Mas apareceu um problema novo: o que acontece quando a aplicação falha entre lotes? Perdes o que estiver no buffer.
Tentativa 3: Redis como buffer
Por isso juntámos o Redis como buffer intermédio. Os eventos chegam, são acrescentados a um Redis Stream (XADD), e um worker em segundo plano consome o stream (XREADGROUP) e despeja lotes no Postgres.
Isto funcionou mesmo, e o Redis continua na arquitetura final. Mas não como armazenamento de dados. Guarda ponteiros para os objetos que chegam, regista o que chegou e o que já foi despejado, e permite-nos fazer batching de forma inteligente sem a aplicação ter de guardar estado. Se a mangueira tem uma pressão de água tão alta que te arrancava a pele, o Redis é a válvula que a transforma em algo manejável.
A ideia-chave foi não usar o Redis como base de dados. É uma camada de coordenação. Os dados de telemetria vão direitinhos para o Postgres. O Redis só diz o que está à espera e o que já foi processado.
A pergunta
Nesta altura tinha um sistema que funcionava. Postgres para armazenamento, Redis para coordenação, escritas COPY em lotes para o débito. Mas não sabia realmente quanto o Postgres aguentava. Aguentaria 10 vezes a carga atual? 100 vezes? Qual é o teto verdadeiro?
Só há uma maneira de descobrir.
Vamos lá só… testar
Decidi fazer o que qualquer pessoa razoável faria: pôr um Postgres a correr num contentor Docker no meu portátil e atirar-lhe quantidades cada vez mais absurdas de telemetria até alguma coisa partir.
A configuração:
- Postgres 16 em Docker com 512MB de shared buffers e 500 ligações no máximo
- pgBouncer em modo de pooling por transação para os testes de pooling de ligações
- Um esquema ao estilo OpenTelemetry: tabelas
spans,metricselogscom atributos JSONB - Um gerador de telemetria que cria traces realistas: ~1 900 eventos por utilizador simulado, com ~3KB cada, distribuídos por ~100 traces com 5 a 12 spans, métricas e logs
- Quatro estratégias de escrita, testadas com números crescentes de utilizadores
- Nos testes extremos (10K+ utilizadores), corremos testes em rajadas de 30 segundos porque, com ~200MB/seg de escritas sustentadas, o meu SSD de 512GB enchia depressa. Mais tempo do que isso e estaria a medir o modo de falha do armazenamento do meu portátil, não o Postgres.
Tudo isto num MacBook Air M2 de 2022 com 24GB de RAM e um SSD de 512GB. Não é um servidor. É um portátil.
As quatro estratégias
INSERT ingénuo: Uma instrução INSERT por evento de telemetria. 50 ligações em simultâneo. A referência do “por favor não faças isto”, e exatamente o que eu fazia em produção às 3 da manhã daquela terça-feira.
INSERT em lotes: Acumular eventos em lotes de 500 linhas e escrevê-los num único INSERT de várias linhas. 20 ligações num pool.
Protocolo COPY: O Postgres tem um protocolo nativo de carga em massa chamado COPY. Faz streaming de dados separados por tabulações diretamente para a tabela, ignorando por completo o parser de SQL. É o que as ferramentas de ETL usam. Lotes de 5 000 linhas encaminhados através do pg-copy-streams.
Pooled + Batched: A mesma estratégia de lotes, mas encaminhada pelo pgBouncer em modo de pooling por transação. Testa se o pooling de ligações acrescenta débito relevante com muita concorrência.
Os números
Com 1 000 utilizadores simulados, cada um a gerar ~1 900 eventos de ~3KB, estamos a escrever aproximadamente 1,8 milhões de linhas de telemetria realista em três tabelas. Payloads JSONB a sério, com nomes de modelos, contagens de tokens, dados de custo, IDs de sessão. O tipo de dados que verias mesmo num pipeline de telemetria de IA em produção.
Eis o que aconteceu.
Escritas por segundo
Com volume baixo, tudo parece mais ou menos igual. Não se nota a diferença entre estratégias quando se escrevem umas poucas milhares de linhas.
Mas sobe a escala e as linhas divergem. A abordagem ingénua estabiliza nas 14 000 escritas/seg e fica por aí. O estrangulamento está na sobrecarga por instrução e na disputa por ligações.
O protocolo COPY chega às 67 000 escritas por segundo. Num portátil. Num contentor Docker. Com max_wal_size a 4GB e as definições de checkpoint por omissão. Com payloads de 3KB, são cerca de 200MB/seg de débito de escrita sustentado. Se tivéssemos afinado os timeouts de checkpoint e a compressão de WAL, provavelmente teríamos ido mais longe.
Um senão do COPY: é tudo ou nada. Se uma linha num lote de 5 000 tiver JSON malformado ou violar uma restrição, o lote inteiro falha. O INSERT em lotes deixa-te usar cláusulas ON CONFLICT para lidar com duplicados e dados sujos com elegância. Em produção, validamos antes do COPY e recorremos ao INSERT em lotes para tudo o que pareça duvidoso.
O INSERT em lotes fica a meio, nas ~34K escritas/seg. Sólido e prático, e não precisas de aprender uma API nova para o fazer.
A estratégia com pooling pelo pgBouncer ficou nas ~43-49K escritas/seg, mais rápida do que o batching direto porque o pgBouncer reutiliza ligações com mais eficiência do que o nosso pool ao nível da aplicação.
Frente a frente com 1 000 utilizadores
Com 1 000 utilizadores (1,8M de eventos), o COPY faz 67K escritas/seg enquanto o ingénuo se arrasta nas 14K. É quase 5 vezes mais rápido para os mesmos dados. Com estes tamanhos de payload, o COPY escreveu 916K eventos em 13,6 segundos. O ingénuo demorou mais de um minuto.
Mas o que me surpreendeu foi isto: nenhuma das estratégias caiu. Esperava que o Postgres começasse a engasgar-se. Erros de ligação, OOM kills, WAL inchado. Nada disso aconteceu. O Postgres limitou-se a aguentar.
Por isso, claro, decidi subir a parada.
Uma verificação da realidade
Sejamos claros: neste momento não temos 10 milhões de utilizadores em simultâneo a enviar 1 900 eventos de telemetria por dia. Se tivéssemos, estaríamos a ganhar dinheiro suficiente para contratar um departamento inteiro para se preocupar com a arquitetura de bases de dados.
Hoje, o nosso volume real de produção é facilmente absorvido pela configuração atual do Postgres sem suar. Então porquê levar a simulação até milhões de utilizadores a gerar centenas de megabytes por segundo?
Estritamente para provar um ponto.
Queria saber o que acontece quando a escolha “aborrecida” finalmente bate na parede. E o que descobri é que, mesmo em escalas absurdas e teóricas em que o Postgres de facto abranda e as consultas degradam, ele não morre sem mais nem menos. Degrada com graciosidade. E, mais importante, quando abranda, tens alavancas normais e aborrecidas para lhe devolver a velocidade.
O teste a sério: escritas E leituras ao mesmo tempo
Eis a questão dos benchmarks de escrita isolados. Estão a mentir-te.
Em produção, não podes pausar as leituras enquanto ingeres dados. O nosso pipeline de análise corre consultas agregadas sobre estes dados enquanto estão a ser escritos. Cálculos de percentis sobre milhões de spans. Mapas de dependências entre serviços. Repartições de taxas de erro. Consultas que obrigam o Postgres a pensar a sério.
Por isso fiz o óbvio: corri escritas COPY em streaming e 8 consultas analíticas em simultâneo, escalando de 1 000 utilizadores até 10 milhões. Nos testes maiores usei janelas de rajada de 30 segundos porque, com estes ritmos de escrita, o SSD de 512GB do meu MacBook encheria fisicamente em menos de 10 minutos.
Com 1 000 utilizadores (escrita completa, 1,8M de eventos), o COPY sustentou 39K escritas/seg enquanto servia, ao mesmo tempo, mais de 7 000 consultas analíticas com um p95 de 9ms. A consulta mais lenta, o JOIN do mapa de dependências, demorou 1,2 segundos.
Com 10 milhões de utilizadores simulados, em rajada de 30 segundos, o Postgres manteve 16 763 escritas/seg enquanto servia consultas analíticas com um p95 abaixo dos 200ms. Em 30 segundos, escreveu 514 536 linhas de telemetria de 3KB. Em todos os pontos de escala, a consulta do mapa de dependências (um self-join sobre a tabela de spans) foi sempre a mais lenta, com um pico de cerca de 1,2 segundos.
O débito de escrita degrada à medida que a tabela cresce e as leituras competem pelo I/O. Era de esperar. Mas o Postgres nunca crashou, nunca ficou sem memória, nunca corrompeu nada. Ficou mais lento e continuou. Todas as consultas devolveram resultados corretos. Todas as escritas foram confirmadas.
E lembra-te: isto correu com apenas 512MB de shared_buffers. A base de dados era bastante maior do que a sua cache. Se lhe tivesse dado 4GB de shared_buffers neste Mac de 24GB, o conjunto de trabalho inteiro viveria em RAM e as consultas seriam muito mais rápidas. Não o fiz de propósito, porque as bases de dados de produção nem sempre podem ter tudo em memória.
Está bem, mas conseguimos corrigir as leituras lentas?
Aquela consulta do mapa de dependências estava a chatear-me. Por isso adicionei índices e repeti o teste sobre ~920K linhas (a telemetria de 500 utilizadores).
Sete índices dirigidos: pesquisas por nome de serviço, filtros por código de estado, joins por ID de trace, relações de spans pai, pesquisas compostas de métricas e distribuições de severidade. Depois ANALYZE para atualizar as estatísticas do planeador de consultas.
O resultado de destaque:
Consulta de erros recentes: 51ms → 1,2ms. Uma aceleração de 43 vezes. Para ser claro, status_code e start_time são colunas de topo no esquema, não estão enterradas no payload JSONB. A coluna JSONB attributes guarda as coisas flexíveis (cabeçalhos HTTP, contagens de tokens, nomes de modelos, dados de custo). Os campos que consultamos com frequência são colunas extraídas, com tipos próprios. O índice parcial em status_code = 2 combinado com o índice start_time DESC permitiu ao Postgres evitar por completo a leitura da tabela inteira. Limita-se a percorrer o índice e a apanhar os 50 primeiros.
Distribuição de severidade dos logs: 65ms → 33ms. O índice composto em (severity, service_name) transforma uma leitura sequencial num index-only scan. 2 vezes mais rápido.
A consulta do mapa de dependências desceu de 456ms para 320ms. Uma melhoria de 1,4 vezes. Melhor, mas não transformadora. Essa consulta faz um self-join sobre centenas de milhares de linhas, e nenhum índice sozinho consegue eliminar o custo fundamental de correlacionar spans pai e filho a esta escala.
Mas tudo bem. Há caminhos de escala claros para quando for preciso.
O plano de escala (para quando for mesmo preciso)
Esta é a parte em que vou argumentar contra a otimização precoce. O que construímos funciona. Aguenta a nossa carga atual com folga, e aguentará 10 vezes mais sem suar. Mas se algum dia chegarmos a processar centenas de milhões de eventos, há dois caminhos claros, ambos ainda com Postgres.
Caminho 1: réplicas de leitura
O passo de escala mais simples. Montas uma réplica com streaming replication e apontas todas as consultas analíticas para ela. As escritas vão para a primária, as leituras para a réplica. A primária pode concentrar-se por inteiro na ingestão, e a réplica pode mastigar JOINs complexos sem afetar o débito de escrita.
É uma mudança para uma tarde de terça-feira. O Google Cloud SQL, o AWS RDS e o Azure suportam todos réplicas de leitura de origem. Acrescentas uma connection string e uma regra de encaminhamento. O teu débito de escrita volta ao patamar das 67K escritas/seg do modo isolado porque já não está a lutar com as leituras, e as tuas consultas analíticas podem demorar o que quiserem na réplica sem ninguém reparar. Em hardware de servidor a sério, com mais núcleos de CPU e armazenamento mais rápido, esse número seria bastante mais alto.
No nosso caso, em que a análise não tem de 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 as tuas tabelas em pedaços. Uma partição por dia, por semana ou por mês. As consultas que filtram por tempo só varrem as partições relevantes em vez da tabela inteira. Aquela consulta do mapa de dependências de 1,2 segundos? Se só olhares para as últimas 24 horas em vez do histórico completo, estás a varrer uma fração dos dados. A consulta passa de segundos a milissegundos.
O particionamento também torna trivial a gestão do ciclo de vida dos dados. Queres largar dados com mais de 90 dias? DROP TABLE spans_2025_q4. Sem vacuum, sem bloat, sem tabela bloqueada. Instantâneo.
A opção de escala gigante
Se alguma vez precisássemos de ir a sério para o disparate, centenas de milhões de utilizadores, milhares de milhões de eventos de telemetria, o Postgres também tem resposta para isso. O Citus é uma extensão do Postgres (totalmente open source, da Microsoft) que distribui os teus dados por vários nós Postgres com sharding por hash. Fazes CREATE EXTENSION citus; e já está. O teu esquema fica igual. As tuas consultas ficam iguais. Só tens mais nós a fazer o trabalho.
Fazes o shard da tabela de spans por trace_id e cada nó trata de uma fatia do conjunto total. Uma consulta a um trace específico bate num nó. Uma consulta agregada espalha-se por todos os nós e junta os resultados. Escala horizontal linear sem sair do ecossistema Postgres. As mesmas ferramentas, a mesma monitorização, a mesma experiência.
Não precisamos disto. Provavelmente não vamos precisar durante muito tempo. Mas o facto de o caminho existir, sem sair do Postgres, é exatamente a razão pela qual começar simples foi a decisão certa.
Não otimizes cedo demais
Quero ser mesmo claro nisto: a arquitetura que temos neste momento não é a mais otimizada para este problema. É a mais apropriada.
Temos um pipeline de dados. Conseguimos fazer consultas complexas enquanto lhe despejamos escritas em cima. Aguenta cargas analíticas em simultâneo. Não cai. Então para que raio precisaria eu de outro produto de base de dados para isto?
Claro que o I/O de disco pode vir a ser um limite um dia. Mas há opções de escala definitivas, e sabemos exatamente quais são porque medimos onde estão os estrangulamentos. Essa é a vantagem de começar simples: quando precisas de escalar, sabes o que escalar.
Podíamos ter começado com ClickHouse. Podíamos ter montado o Citus desde o primeiro dia. Podíamos ter construído uma arquitetura Lambda com Kafka e um processador de streams e uma camada de serving e uma camada de batch e uma… já percebeste.
Mas toda essa complexidade tem um custo. Cada sistema adicional é mais uma coisa para monitorizar, mais uma coisa que pode falhar às 3 da manhã, mais uma coisa que a tua equipa tem de perceber. Se não tens a escala que o justifique, estás a pagar 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 para ingestão em lotes. Essa é a nossa arquitetura. Aguenta a nossa carga atual. Aguentará 10 vezes a nossa carga atual. E quando não aguentar 100 vezes, saberemos exatamente onde estão os estrangulamentos porque os medimos.
A boa notícia é que não precisamos de que este serviço seja rápido. Os dados de telemetria são processados em lotes para análise, não servidos em tempo real aos utilizadores. Uma consulta que demora 4 segundos em vez de 40 milissegundos é perfeitamente aceitável para o nosso caso.
Mas se alguma vez precisássemos que fosse rápido… bem, já viste os números. Há muita margem.
O que aprendi de facto
- Nunca uses INSERTs individuais para escritas de alto débito. Lotes ou COPY. Sempre.
- O protocolo COPY existe por uma razão. Não serve só para cargas iniciais de dados. É uma estratégia legítima de ingestão em produção.
- O pooling de ligações não é opcional à escala. pgBouncer em modo de transação. Sem desculpas.
- Testa escritas E leituras em conjunto. Os benchmarks só de escrita enganam. O perfil de desempenho real é o que acontece quando a tua base de dados faz as duas coisas ao mesmo tempo.
- Os índices importam, mas não para tudo. Um índice parcial bem colocado pode dar-te uma aceleração de 43 vezes em consultas dirigidas. Mas as consultas agregadas sobre milhões de linhas vão ser lentas aconteça o que acontecer. É aí que recorres a réplicas e particionamento.
- Não otimizes cedo demais. Começa com a arquitetura mais simples que funcione. Mede onde estão os estrangulamentos. Escala as partes que precisam de escalar. Nem tudo, nem tudo de uma vez.
- O Postgres aguenta mais do que as pessoas lhe reconhecem. 67K escritas/seg de eventos de telemetria de 3KB. 39K escritas/seg enquanto serve consultas analíticas ao mesmo tempo. Escalado para 10M de utilizadores simulados e continuou sem crashar. Num portátil. Em Docker.
- Testa os teus pressupostos. Passei semanas a ler posts de blog com comparações de data warehouses. Podia ter passado uma tarde com o Docker e um script e ter números reais. Os números contaram uma história melhor.
Ah, e mais uma coisa. Passámos este post inteiro a pôr o Postgres num inferno absoluto. Milhões de linhas, leituras em simultâneo, self-joins sobre spans gigantescos. E nem sequer pensámos no Redis que está à frente dele. O Redis que coordena tudo isto, que trata da mangueira de telemetria a chegar, que segue os ponteiros, que gere o estado dos lotes. Nem pestanejou. Não o testámos porque não havia nada para testar. Funcionou e pronto.
A maioria dos problemas resolve-se mesmo com um servidor web, um Redis e um Postgres.
O código do benchmark está escrito em Go e vive em github.com/willhackett/bench-postgres. docker-compose up -d, go run . -mode=full para as estratégias de escrita, go run . -mode=chaos para o teste de caos de escritas e leituras combinadas e go run . -mode=optimize para a comparação de indexação.