← Tous les articles

L'ennuyeux gagne : 67k événements de télémétrie par seconde dans Postgres

Sur cette page

On collecte énormément de télémétrie. Des prompts, des métriques d’usage, des données de session, des informations venues de chaque assistant de code IA et de chaque application de bureau que les ingénieurs de nos clients utilisent. Tout ça afflue de milliers d’utilisateurs vers notre plateforme d’analyse.

Chez Flowstate, on utilise OpenTelemetry pour mesurer les dépenses d’IA. Chaque appel d’API, chaque session de code, chaque invocation de modèle, rattachés aux équipes, aux projets et aux centres de coûts. Le volume est important et ça ne s’arrête jamais.

Du coup, il faut stocker beaucoup de données. Pas des données de « tableau de bord analytique ». Pas des données de « rapport mensuel ». Chaque span, chaque métrique, chaque ligne de log de chaque utilisateur chez chaque client, conservés pour que notre pipeline d’analyse puisse les mâcher plus tard.

Tout le monde a le même premier réflexe : « Il te faut un data warehouse. »

J’ai essayé. J’ai vraiment essayé.

Tour du marché

Précisons d’abord que tous les produits que j’ai regardés sont excellents pour ce pour quoi ils sont conçus. Ils visent des organisations avec de grosses équipes data et des charges analytiques complexes. Notre problème était plus simple : écrire beaucoup de télémétrie vite, l’interroger plus tard. Pour ça, la plupart de ces solutions dépassaient nos besoins.

Il y avait aussi une contrainte pratique qui me trottait dans la tête. Je dois pouvoir développer dans un avion. Pas que je le fasse vraiment, mais c’est le principe. Si un composant de ma stack exige une connexion internet pour fonctionner, je ne peux pas construire, tester et itérer en local. Je ne peux pas le lancer dans Docker et lui balancer des données un samedi matin. Beaucoup de ces data warehouses me semblaient vouloir résoudre des problèmes qu’un Postgres bien réglé sait gérer.

BigQuery a été la première étape. Excellent pour interroger à grande échelle, moins adapté aux écritures continues à gros volume. L’API d’insertion en streaming facture à la ligne, et la note grimpe vite à notre volume. Le profil de latence ne convient pas non plus à une ingestion quasi temps réel. Un moteur d’analyse phénoménal, mais pas le bon outil pour un pipeline de télémétrie où l’écriture domine.

Snowflake est une plateforme de données puissante, mais le modèle de prix est complexe (crédits de calcul, stockage, transfert de données) et, pour notre besoin, c’était beaucoup plus d’infrastructure que le problème n’en demandait. Quand l’alternative est un Postgres bien réglé, la comparaison des coûts est brutale.

Databricks est impressionnant si tu as une équipe de data engineering dédiée. L’architecture lakehouse et l’intégration Spark sont réellement puissantes. Mais on n’a pas douze data engineers, et ajouter une plateforme de cette complexité pour ce qui revient à « écrire des lignes, lire des lignes », c’était sortir un drone Predator pour une bagarre au couteau.

Amazon Redshift, Azure Synapse et les autres offres analytiques managées sont toutes de bons produits, mais chacune apporte une charge d’exploitation qui ne collait pas avec notre stade. Un système de plus à surveiller, un jeu d’identifiants de plus, une relation fournisseur de plus.

ClickHouse était l’option la plus séduisante. Orienté colonnes, conçu pour les charges analytiques à forte écriture, open source et vraiment rapide. Je l’ai beaucoup aimé. Si je l’ai écarté, ce n’est pas à cause de la technologie, mais du déploiement. Le faire tourner de façon fiable dans un environnement managé me demandait plus de friction que je ne voulais. ClickHouse Cloud existe, mais c’est encore une dépendance à un fournisseur, alors que j’ai déjà du Postgres managé sur Google Cloud SQL. Et on peut faire de l’analytique dans Postgres de toute façon, donc ClickHouse ressemblait à des étapes en plus.

Ce qui m’a ramené à la base de données que j’avais déjà.

J’ai construit des trucs assez fous sur Postgres au fil des ans, et la question revient toujours : « Tu es sûr qu’une seule base suffit ? » Je peux te dire par expérience qu’une seule base encaisse sévèrement. Des gens font tourner Postgres à l’échelle du pétaoctet. C’est faisable, avec difficulté, mais c’est faisable. Voyons donc ce qu’on peut tirer de notre modeste cas d’usage.

J’ai déjà Postgres. Il marche déjà. Cherchons où il cesse de marcher.

Les échecs

Voici un résumé chronologique de ce que j’ai appris à la dure.

Essai 1 : écrire directement dans la base de production

La première version était aussi bête qu’on l’imagine. À chaque événement de télémétrie reçu, on envoyait un INSERT à l’instance Postgres de production. Une ligne à la fois. Une par une. Comme on poste des lettres.

Ça marchait à faible volume. Ça marchait à volume moyen. Ça a cessé de marcher à 3 h du matin un mardi, quand le service proxy a pris un pic de trafic et qu’on s’est retrouvés à 10 000 INSERT individuels par seconde sur une base qui devait aussi servir l’application.

Je n’arrivais littéralement pas à obtenir un verrou en écriture pour lancer une migration. Le pool de connexions était saturé, chaque connexion disponible était bloquée sur un INSERT, et l’outil de migration attendait un verrou qui ne viendrait jamais. J’ai fini par tuer les connexions à la main pour faire passer la migration. À 3 h du matin. Un mardi.

Essai 2 : le batching (ça se réchauffe)

La correction évidente : arrêter d’écrire une ligne à la fois. Mettre les événements en tampon en mémoire et les vider dans la base par lots de 500 à 1 000 lignes avec des INSERT multi-lignes.

C’était nettement mieux. Au lieu de 10 000 transactions, tu en as 10 à 20 plus grosses. La base respire. Le pool de connexions se vide. Les migrations passent.

Mais un nouveau problème est apparu : que se passe-t-il si l’application plante entre deux lots ? Tu perds ce qu’il y a dans le tampon.

Essai 3 : Redis comme tampon

On a donc ajouté Redis comme tampon intermédiaire. Les événements arrivent, sont ajoutés à un Redis Stream (XADD), et un worker en arrière-plan consomme le flux (XREADGROUP) et vide des lots dans Postgres.

Ça a vraiment marché, et Redis est toujours dans l’architecture finale. Mais pas comme magasin de données. Il stocke des pointeurs vers les objets entrants, garde la trace de ce qui est arrivé et de ce qui a été écrit, et nous laisse faire du batching intelligent sans que l’application ait à garder un état. Si le débit est assez violent pour t’arracher la peau, Redis est la vanne qui le rend gérable.

L’idée clé a été de ne pas utiliser Redis comme base de données. C’est une couche de coordination. Les données de télémétrie elles-mêmes vont directement dans Postgres. Redis nous dit seulement ce qui attend et ce qui a été traité.

La question

À ce stade, j’avais un système qui marchait. Postgres pour le stockage, Redis pour la coordination, des écritures COPY par lots pour le débit. Mais je ne savais pas vraiment combien Postgres pouvait encaisser. Pouvait-il tenir 10 fois la charge actuelle ? 100 fois ? Quel est le vrai plafond ?

Il n’y a qu’une façon de le savoir.

Testons, tout simplement

J’ai fait ce que ferait n’importe quelle personne raisonnable : lancer Postgres dans un conteneur Docker sur mon portable et lui balancer des quantités de télémétrie de plus en plus absurdes jusqu’à ce que quelque chose casse.

Le montage :

  • Postgres 16 dans Docker, avec 512 Mo de shared buffers et 500 connexions max
  • pgBouncer en mode pooling par transaction pour les tests de pool de connexions
  • Un schéma de type OpenTelemetry : des tables spans, metrics et logs avec des attributs JSONB
  • Un générateur de télémétrie qui crée des traces réalistes : environ 1 900 événements par utilisateur simulé, d’environ 3 Ko chacun, répartis sur environ 100 traces de 5 à 12 spans, avec des métriques et des logs
  • Quatre stratégies d’écriture, testées sur un nombre croissant d’utilisateurs
  • Pour les tests extrêmes (10K+ utilisateurs), on lance des rafales de 30 secondes, parce qu’à environ 200 Mo/s d’écriture soutenue, mon SSD de 512 Go se remplirait vite. Plus longtemps, et je mesurerais la façon dont mon stockage tombe en panne, pas Postgres.

Le tout sur un MacBook Air M2 de 2022 avec 24 Go de RAM et un SSD de 512 Go. Pas un serveur. Un portable.

Les quatre stratégies

INSERT naïf : un INSERT par événement de télémétrie. 50 connexions simultanées. La référence « s’il vous plaît, ne faites pas ça », et exactement ce que je faisais en production à 3 h du matin ce mardi-là.

INSERT par lots : on accumule les événements par lots de 500 lignes et on les écrit dans un seul INSERT multi-lignes. 20 connexions dans un pool.

Protocole COPY : Postgres a un protocole intégré de chargement en masse, COPY. Il envoie des données séparées par des tabulations directement dans la table, sans passer par le parseur SQL. C’est ce que les outils d’ETL utilisent. Des lots de 5 000 lignes passés par pg-copy-streams.

Pool + lots : même stratégie de lots, mais routée par pgBouncer en mode pooling par transaction. Ça teste si le pooling de connexions ajoute un vrai débit à forte concurrence.

Les chiffres

Avec 1 000 utilisateurs simulés, chacun générant environ 1 900 événements d’environ 3 Ko, on écrit à peu près 1,8 million de lignes de télémétrie réaliste réparties sur trois tables. De vrais payloads JSONB avec des noms de modèles, des nombres de tokens, des données de coût, des identifiants de session. Le genre de données qu’on verrait dans un vrai pipeline de télémétrie IA en production.

Voici ce qui s’est passé.

Écritures par seconde

À faible volume, tout se ressemble. Tu ne vois aucune différence entre les stratégies quand tu écris quelques milliers de lignes.

Mais monte en charge et les courbes s’écartent. L’approche naïve plafonne à environ 14 000 écritures/s et n’en bouge plus. Tu es freiné par le coût de chaque requête et la contention sur les connexions.

Le protocole COPY atteint 67 000 écritures par seconde. Sur un portable. Dans un conteneur Docker. Avec max_wal_size à 4 Go et les réglages de checkpoint par défaut. Avec des payloads de 3 Ko, ça fait à peu près 200 Mo/s de débit d’écriture soutenu. Si on avait réglé les délais de checkpoint et la compression du WAL, on aurait probablement pu monter plus haut.

Une réserve sur COPY : c’est tout ou rien. Si une seule ligne d’un lot de 5 000 contient un JSON mal formé ou viole une contrainte, tout le lot échoue. L’INSERT par lots permet d’utiliser des clauses ON CONFLICT pour gérer proprement les doublons et les données sales. En production, on prévalide avant COPY et on retombe sur l’INSERT par lots pour tout ce qui paraît louche.

L’INSERT par lots se place au milieu, à environ 34K écritures/s. Solide et pratique, et pas besoin d’apprendre une nouvelle API.

La stratégie avec pool via pgBouncer atteint environ 43 à 49K écritures/s, plus vite que le batching brut parce que pgBouncer réutilise les connexions plus efficacement que notre pool côté application.

Face à face à 1 000 utilisateurs

À 1 000 utilisateurs (1,8 M d’événements), COPY tient 67K écritures/s quand l’approche naïve reste bloquée à 14K. C’est près de 5 fois plus rapide pour les mêmes données. Avec ces tailles de payload, COPY a écrit 916K événements en 13,6 secondes. L’approche naïve a mis plus d’une minute.

Mais voilà ce qui m’a surpris : aucune des stratégies n’a flanché. Je m’attendais à ce que Postgres s’étouffe. Erreurs de connexion, OOM kills, gonflement du WAL. Rien de tout ça n’est arrivé. Postgres a juste encaissé.

Alors, naturellement, j’ai poussé le curseur.

Un rappel à la réalité

Soyons clairs : on n’a pas actuellement 10 millions d’utilisateurs simultanés qui envoient 1 900 événements de télémétrie par jour. Si c’était le cas, on gagnerait assez d’argent pour embaucher un département entier chargé de s’inquiéter de l’architecture des bases de données.

Aujourd’hui, notre volume réel de production est géré sans effort par notre Postgres actuel. Alors pourquoi pousser la simulation jusqu’à des millions d’utilisateurs qui produisent des centaines de mégaoctets par seconde ?

Uniquement pour prouver quelque chose.

Je voulais savoir ce qui se passe quand le choix « ennuyeux » finit par toucher le mur. Et ce que j’ai trouvé, c’est que même à des échelles absurdes, théoriques, où Postgres ralentit et où les requêtes se dégradent, il ne meurt pas d’un coup. Il se dégrade en douceur. Et surtout, quand il ralentit, tu as des leviers standard, ennuyeux, pour retrouver la vitesse.

Le vrai test : écrire ET lire en même temps

Voilà le problème des benchmarks d’écriture isolés. Ils te mentent.

En production, tu ne peux pas mettre les lectures en pause pendant que tu ingères des données. Notre pipeline d’analyse lance des requêtes d’agrégation sur ces données pendant qu’elles sont écrites. Des calculs de percentiles sur des millions de spans. Des cartes de dépendances entre services. Des ventilations de taux d’erreur. Des requêtes qui font réfléchir Postgres.

J’ai donc fait l’évidence : lancer des écritures COPY en streaming et 8 requêtes analytiques en même temps, de 1 000 utilisateurs jusqu’à 10 millions. Pour les plus gros tests, j’ai utilisé des fenêtres de rafale de 30 secondes, parce qu’à ces débits d’écriture, le SSD de 512 Go de mon MacBook serait physiquement plein en moins de 10 minutes.

À 1 000 utilisateurs (écriture complète, 1,8 M d’événements), COPY a soutenu 39K écritures/s tout en servant plus de 7 000 requêtes analytiques avec un p95 de 9 ms. La requête la plus lente, le JOIN de la carte de dépendances, a pris 1,2 seconde.

À 10 millions d’utilisateurs simulés, en rafale de 30 secondes, Postgres a maintenu 16 763 écritures/s tout en servant des requêtes analytiques avec un p95 sous 200 ms. En 30 secondes, il a écrit 514 536 lignes de télémétrie de 3 Ko. À toutes les échelles, la requête de carte de dépendances (une auto-jointure sur la table spans) a été la plus lente, avec un pic autour de 1,2 seconde.

Le débit d’écriture se dégrade à mesure que la table grossit et que les lectures se disputent les I/O. C’est attendu. Mais Postgres n’a jamais planté, n’a jamais subi d’OOM, n’a rien corrompu. Il a ralenti et il a continué. Chaque requête a renvoyé des résultats corrects. Chaque écriture a été validée.

Et rappelle-toi : tout ça tournait avec seulement 512 Mo de shared_buffers. La base était bien plus grosse que son cache. Si je lui avais donné 4 Go de shared_buffers sur ce Mac de 24 Go, tout le jeu de données de travail aurait tenu en RAM et ces requêtes auraient été bien plus rapides. Je ne l’ai pas fait exprès, parce que les bases de production ne peuvent pas toujours tout garder en mémoire.

Les lectures lentes se corrigent

Cette requête de carte de dépendances me turlupinait. J’ai donc ajouté des index et relancé le test sur environ 920K lignes (la télémétrie de 500 utilisateurs).

Sept index ciblés : recherches par nom de service, filtres sur le code de statut, jointures par ID de trace, relations de span parent, recherches composites sur les métriques et distributions de sévérité. Puis ANALYZE pour mettre à jour les statistiques du planificateur de requêtes.

Le résultat qui se détache :

Requête des erreurs récentes : 51 ms → 1,2 ms. Une accélération de 43 fois. Précisons que status_code et start_time sont des colonnes de premier niveau dans le schéma, pas enfouies dans le payload JSONB. La colonne JSONB attributes stocke le contenu flexible (en-têtes HTTP, nombres de tokens, noms de modèles, données de coût). Les champs qu’on interroge souvent sont des colonnes extraites, correctement typées. L’index partiel sur status_code = 2, combiné à l’index sur start_time DESC, a permis à Postgres d’éviter complètement le scan de la table. Il parcourt l’index et prend les 50 premières lignes.

Distribution de la sévérité des logs : 65 ms → 33 ms. L’index composite sur (severity, service_name) transforme un scan séquentiel en scan d’index seul. 2 fois plus rapide.

La requête de carte de dépendances est passée de 456 ms à 320 ms. Une amélioration de 1,4 fois. Mieux, mais pas transformatrice. Cette requête fait une auto-jointure sur des centaines de milliers de lignes, et aucun index ne peut éliminer le coût de fond qu’il y a à corréler des spans parents et enfants à cette échelle.

Mais ce n’est pas grave. Il existe des chemins clairs pour monter en charge le jour où il le faut.

La feuille de route pour monter en charge (quand tu en as vraiment besoin)

C’est ici que je plaide contre l’optimisation prématurée. Ce qu’on a construit marche. Ça gère notre charge actuelle confortablement, et ça en gérerait 10 fois plus sans effort. Mais si on en vient un jour à traiter des centaines de millions d’événements, il y a deux chemins clairs, tous deux toujours sous Postgres.

Chemin 1 : les réplicas de lecture

Le levier le plus simple. Tu mets en place un réplica en streaming et tu envoies dessus toutes les requêtes analytiques. Les écritures vont au primaire, les lectures au réplica. Le primaire peut se consacrer entièrement à l’ingestion, et le réplica peut mâcher des JOIN complexes sans toucher au débit d’écriture.

C’est un changement d’un mardi après-midi. Google Cloud SQL, AWS RDS et Azure gèrent tous les réplicas de lecture nativement. Tu ajoutes une chaîne de connexion et une règle de routage. Ton débit d’écriture revient à la zone des 67K écritures/s en autonome, parce qu’il ne se bat plus contre les lectures, et tes requêtes analytiques peuvent prendre tout le temps qu’elles veulent sur le réplica sans que personne le remarque. Sur du vrai matériel de serveur, avec plus de cœurs CPU et un stockage plus rapide, ce chiffre serait nettement plus élevé.

Pour notre cas, où l’analyse n’a pas besoin d’être en temps réel, un réplica avec quelques secondes de retard de réplication convient parfaitement.

Chemin 2 : le partitionnement des tables

Le partitionnement par le temps découpe tes tables en morceaux. Une partition par jour, par semaine ou par mois. Les requêtes qui filtrent par date ne scannent que les partitions concernées au lieu de la table entière. Cette requête de carte de dépendances à 1,2 seconde ? Si tu ne regardes que les dernières 24 heures au lieu de tout l’historique, tu scannes une fraction des données. La requête passe de secondes à millisecondes.

Le partitionnement rend aussi la gestion du cycle de vie des données triviale. Tu veux supprimer les données de plus de 90 jours ? DROP TABLE spans_2025_q4. Pas de vacuum, pas de bloat, pas de table verrouillée. Instantané.

L’option méga-échelle

Si on avait un jour besoin de passer vraiment à l’excès, avec des centaines de millions d’utilisateurs et des milliards d’événements de télémétrie, Postgres a une réponse aussi. Citus est une extension Postgres (entièrement ouverte par Microsoft) qui répartit tes données sur plusieurs nœuds Postgres par sharding par hachage. Tu fais CREATE EXTENSION citus; et c’est parti. Ton schéma reste le même. Tes requêtes restent les mêmes. Tu as juste plus de nœuds pour faire le travail.

Tu shardes la table spans par trace_id et chaque nœud gère une tranche de l’ensemble. Une requête sur une trace précise tape un seul nœud. Une requête d’agrégation se déploie sur tous les nœuds et fusionne les résultats. Une montée en charge horizontale linéaire, tout en restant dans l’écosystème Postgres. Mêmes outils, même supervision, même expertise.

On n’en a pas besoin. On n’en aura probablement pas besoin avant très longtemps. Mais le fait que ce chemin existe, sans quitter Postgres, c’est exactement ce qui prouve que partir simple était le bon choix.

N’optimise pas trop tôt

Je veux être très clair : l’architecture qu’on fait tourner aujourd’hui n’est pas l’architecture la plus optimale pour ce problème. C’est la plus appropriée.

On a un pipeline de données. On peut lancer des requêtes complexes pendant qu’on le spamme d’écritures. Il gère des charges analytiques concurrentes. Il ne s’effondre pas. Alors pourquoi aurais-je besoin d’un autre produit de base de données ?

Bien sûr, les I/O disque pourraient devenir une limite un jour. Mais il existe des options de montée en charge bien définies, et on sait exactement lesquelles, parce qu’on a mesuré où se trouvent les goulots. C’est l’avantage de partir simple : quand tu dois monter en charge, tu sais quoi faire monter.

On aurait pu commencer avec ClickHouse. On aurait pu installer Citus dès le premier jour. On aurait pu bâtir une architecture Lambda avec Kafka, un stream processor, une couche de service, une couche batch et une… tu vois l’idée.

Mais toute cette complexité a un coût. Chaque système en plus est une chose de plus à surveiller, une chose de plus qui peut tomber en panne à 3 h du matin, une chose de plus que ton équipe doit comprendre. Si tu n’as pas l’échelle pour la justifier, tu paies la taxe de la complexité sans toucher le bénéfice de la scalabilité.

Postgres sur Cloud SQL, avec Redis comme couche de coordination, et le protocole COPY pour l’ingestion par lots. Voilà notre architecture. Elle gère notre charge actuelle. Elle gérera 10 fois cette charge. Et le jour où elle ne tiendra pas 100 fois, on saura exactement où sont les goulots, parce qu’on les a mesurés.

La bonne nouvelle, c’est qu’on n’a pas besoin que ce service soit rapide. Les données de télémétrie sont traitées par lots pour l’analyse, elles ne sont pas servies en temps réel aux utilisateurs. Une requête qui prend 4 secondes au lieu de 40 millisecondes ne gêne absolument pas notre cas d’usage.

Mais si on avait un jour besoin qu’il soit rapide… eh bien, tu as vu les chiffres. Il y a de la marge.

Ce que j’ai vraiment appris

  • N’utilise jamais d’INSERT individuels pour des écritures à fort débit. Des lots ou COPY. Toujours.
  • Le protocole COPY existe pour une raison. Il ne sert pas qu’aux chargements initiaux. C’est une vraie stratégie d’ingestion en production.
  • Le pooling de connexions n’est pas optionnel à grande échelle. pgBouncer en mode transaction. Pas d’excuses.
  • Teste les écritures ET les lectures ensemble. Les benchmarks d’écriture seule sont trompeurs. Le vrai profil de performance, c’est ce qui arrive quand ta base fait les deux à la fois.
  • Les index comptent, mais pas pour tout. Un index partiel bien placé peut donner une accélération de 43 fois sur des requêtes ciblées. Mais les requêtes d’agrégation sur des millions de lignes resteront lentes quoi que tu fasses. C’est là qu’on sort les réplicas et le partitionnement.
  • N’optimise pas trop tôt. Pars de l’architecture la plus simple qui marche. Mesure où sont les goulots. Fais monter les parties qui en ont besoin. Pas tout, pas tout d’un coup.
  • Postgres encaisse plus que ce qu’on lui reconnaît. 67K écritures/s de télémétrie de 3 Ko. 39K écritures/s tout en servant des requêtes analytiques. Poussé jusqu’à 10 M d’utilisateurs simulés, et il n’a toujours pas planté. Sur un portable. Dans Docker.
  • Teste tes hypothèses. J’ai passé des semaines à lire des articles comparant des data warehouses. J’aurais pu passer un après-midi avec Docker et un script et avoir de vrais chiffres. Les chiffres racontaient une meilleure histoire.

Ah, et une dernière chose. On a passé tout ce billet à faire vivre un enfer à Postgres. Des millions de lignes, des lectures concurrentes, des auto-jointures sur d’énormes tables de spans. Et on n’a même pas pensé au Redis posté devant. Le Redis qui coordonne tout ça, encaisse le torrent de télémétrie entrante, suit les pointeurs, gère l’état des lots. Il n’a même pas bronché. On ne l’a pas benchmarké parce qu’il n’y avait rien à benchmarker. Il a juste marché.

La plupart des problèmes se règlent vraiment avec un serveur web, un Redis et un Postgres.

Le code du benchmark est écrit en Go et se trouve sur github.com/willhackett/bench-postgres. docker-compose up -d, go run . -mode=full pour les stratégies d’écriture, go run . -mode=chaos pour le test de chaos combinant écriture et lecture, et go run . -mode=optimize pour la comparaison d’indexation.